FazBrowse GitHub Viewer | Trending |
URL:
| Home
Tools: [Download Repo ZIP]   [Original HTTPS Page]

fix: Read Kafka and Kinesis sources without a batch source from proto by LuisFigueroaG · Pull Request #6949 · feast-dev/feast · GitHub

Repository navigation

Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension .py  (2) All 1 file type selected
Viewed files
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Unified
Split
Hide whitespace
Diff view
Unified
Split
Hide whitespace
4 changes: 2 additions & 2 deletions sdk/python/feast/data_source.py
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters. Learn more about bidirectional Unicode characters
Original file line number Diff line number Diff line change
Expand Up @@ -570,7 +570,7 @@ def from_proto(data_source: DataSourceProto):
owner=data_source.owner,
batch_source=(
DataSource.from_proto(data_source.batch_source)
if data_source.batch_source
if data_source.HasField("batch_source")
else None
),
)
Expand Down Expand Up @@ -757,7 +757,7 @@ def from_proto(data_source: DataSourceProto):
owner=data_source.owner,
batch_source=(
DataSource.from_proto(data_source.batch_source)
if data_source.batch_source
if data_source.HasField("batch_source")
else None
),
)
Expand Down
38 changes: 38 additions & 0 deletions sdk/python/tests/unit/test_data_sources.py
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters. Learn more about bidirectional Unicode characters
Original file line number Diff line number Diff line change
Expand Up @@ -379,3 +379,41 @@ def kafka_source(watermark):
KafkaSource.from_proto(unset_proto).kafka_options.watermark_delay_threshold
is None
)


@pytest.mark.parametrize(
"stream_source",
[
KafkaSource(
name="test_source",
kafka_bootstrap_servers="test_servers",
message_format=ProtoFormat("class_path"),
topic="test_topic",
timestamp_field="event_timestamp",
),
KinesisSource(
name="test_source",
region="test_region",
record_format=ProtoFormat("class_path"),
stream_name="test_stream",
timestamp_field="event_timestamp",
),
],
ids=["kafka", "kinesis"],
)
def test_stream_source_without_batch_source_round_trips(stream_source):
"""``batch_source`` is optional on Kafka and Kinesis sources.

``from_proto`` checked ``if data_source.batch_source``, but an unset proto
sub-message is still truthy, so the empty message was parsed as a data
source and failed with "Could not identify the source type being added."
That made these sources impossible to read back from the registry.
"""
proto = stream_source.to_proto()
assert not proto.HasField("batch_source")

restored = DataSource.from_proto(proto)

assert isinstance(restored, type(stream_source))
assert restored.batch_source is None
assert restored.name == stream_source.name

Back | FazBrowse Home | New Git URL