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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion VERSION
Original file line number Diff line number Diff line change
@@ -1 +1 @@
3.0.2
3.1.0
8 changes: 8 additions & 0 deletions src/glassflow/etl/models/source.py
Original file line number Diff line number Diff line change
Expand Up @@ -85,6 +85,7 @@ class TopicConfig(BaseModel):
name: str
event_schema: Schema = Field(alias="schema")
deduplication: Optional[DeduplicationConfig] = Field(default=DeduplicationConfig())
replicas: Optional[int] = Field(default=1)

@field_validator("deduplication")
@classmethod
Expand Down Expand Up @@ -119,6 +120,13 @@ def validate_deduplication_id_field(

return v

@field_validator("replicas")
@classmethod
def validate_replicas(cls, v: int) -> int:
if v < 1:
raise ValueError("Replicas must be at least 1")
return v


class KafkaConnectionParams(BaseModel):
brokers: List[str]
Expand Down
2 changes: 2 additions & 0 deletions tests/data/pipeline_configs.py
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ def get_valid_pipeline_config() -> dict:
{
"consumer_group_initial_offset": "earliest",
"name": "user_logins",
"replicas": 3,
"schema": {
"type": "json",
"fields": [
Expand Down Expand Up @@ -52,6 +53,7 @@ def get_valid_pipeline_config() -> dict:
{
"consumer_group_initial_offset": "earliest",
"name": "orders",
"replicas": 1,
"schema": {
"type": "json",
"fields": [
Expand Down
31 changes: 31 additions & 0 deletions tests/test_models/test_topic.py
Original file line number Diff line number Diff line change
Expand Up @@ -71,3 +71,34 @@ def test_topic_config_deduplication_id_field_validation(self):
),
)
assert config.deduplication.enabled is False

def test_topic_config_replicas_validation(self):
"""Test TopicConfig validation for replicas."""
config = models.TopicConfig(
name="test-topic",
consumer_group_initial_offset=models.ConsumerGroupOffset.EARLIEST,
schema=models.Schema(
type=models.SchemaType.JSON,
fields=[
models.SchemaField(name="name", type=models.KafkaDataType.STRING),
],
),
replicas=3,
)
assert config.replicas == 3

with pytest.raises(ValueError) as exc_info:
models.TopicConfig(
name="test-topic",
consumer_group_initial_offset=models.ConsumerGroupOffset.EARLIEST,
schema=models.Schema(
type=models.SchemaType.JSON,
fields=[
models.SchemaField(
name="name", type=models.KafkaDataType.STRING
),
],
),
replicas=0,
)
assert "Replicas must be at least 1" in str(exc_info.value)