Skip to content
Open
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
6 changes: 6 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,12 @@ All notable changes to this project will be documented in this file.
The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/),
and this project follows [Semantic Versioning](https://semver.org/spec/v2.0.0.html).

## [Unreleased]

### Added

- `define_sink_pipe` now supports GCS as a sink/export destination, matching the Forward CLI's `gcs_hmac` export service. Add a `GCSConnectionDefinition` as the sink's `connection` and the same `bucket_uri`/`file_template`/`format`/`schedule`/`strategy`/`compression` options already used for S3 sinks; the datafile emitter writes an explicit `EXPORT_SERVICE gcs_hmac` directive (since GCS and S3 sinks otherwise share identical `EXPORT_*` directives) and the migration parser/emitter round-trip it accordingly.

## [0.4.0] - 2026-06-29

### Added
Expand Down
26 changes: 25 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -772,7 +772,7 @@ manual_report = define_copy_pipe(
### Sink Pipes

Use sink pipes to publish query results to external systems.
The SDK supports Kafka and S3 sinks.
The SDK supports Kafka, S3, and GCS sinks.

```python
from tinybird_sdk import define_sink_pipe, node
Expand Down Expand Up @@ -820,6 +820,30 @@ s3_events_sink = define_sink_pipe(
],
},
)

# GCS sink
gcs_events_sink = define_sink_pipe(
"gcs_events_sink",
{
"sink": {
"connection": landing_gcs,
"bucket_uri": "gs://my-bucket/exports/",
"file_template": "events_{date}",
"format": "csv",
"schedule": "@once",
"strategy": "create_new",
"compression": "gzip",
},
"nodes": [
node(
{
"name": "export",
"sql": "SELECT timestamp, session_id FROM gcs_landing",
}
)
],
},
)
```

### Static Tokens
Expand Down
13 changes: 13 additions & 0 deletions src/tinybird_sdk/generator/pipe.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
from ..schema.pipe import (
CopyConfig,
EndpointConfig,
GCSSinkConfig,
KafkaSinkConfig,
MaterializedConfig,
PipeDefinition,
Expand Down Expand Up @@ -83,6 +84,18 @@ def _generate_sink(config: SinkConfig) -> str:
parts.append(f"EXPORT_STRATEGY {config.strategy}")
if config.compression:
parts.append(f"EXPORT_COMPRESSION {config.compression}")
elif isinstance(config, GCSSinkConfig):
# GCS and S3 sinks share the same EXPORT_* directive shape, so EXPORT_SERVICE
# must be emitted explicitly to disambiguate gcs_hmac from s3 on re-parse.
parts.append("EXPORT_SERVICE gcs_hmac")
parts.append(f"EXPORT_BUCKET_URI {config.bucket_uri}")
parts.append(f"EXPORT_FILE_TEMPLATE {config.file_template}")
parts.append(f"EXPORT_SCHEDULE {config.schedule}")
parts.append(f"EXPORT_FORMAT {config.format}")
if config.strategy:
parts.append(f"EXPORT_STRATEGY {config.strategy}")
if config.compression:
parts.append(f"EXPORT_COMPRESSION {config.compression}")

return "\n".join(parts)

Expand Down
1 change: 1 addition & 0 deletions src/tinybird_sdk/migrate/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@
PipeModel,
SinkKafkaModel,
SinkS3Model,
SinkGCSModel,
SinkModel,
KafkaConnectionModel,
S3ConnectionModel,
Expand Down
3 changes: 2 additions & 1 deletion src/tinybird_sdk/migrate/emit_ts.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@
ParsedResource,
PipeModel,
S3ConnectionModel,
SinkGCSModel,
SinkKafkaModel,
SinkS3Model,
)
Expand Down Expand Up @@ -427,7 +428,7 @@ def _emit_pipe(pipe: PipeModel) -> str:
if isinstance(pipe.sink, SinkKafkaModel):
lines.append(f" 'topic': {_escape_string(pipe.sink.topic)},")
lines.append(f" 'schedule': {_escape_string(pipe.sink.schedule)},")
elif isinstance(pipe.sink, SinkS3Model):
elif isinstance(pipe.sink, (SinkS3Model, SinkGCSModel)):
lines.append(f" 'bucket_uri': {_escape_string(pipe.sink.bucket_uri)},")
lines.append(f" 'file_template': {_escape_string(pipe.sink.file_template)},")
lines.append(f" 'schedule': {_escape_string(pipe.sink.schedule)},")
Expand Down
53 changes: 38 additions & 15 deletions src/tinybird_sdk/migrate/parse_pipe.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
PipeTokenModel,
PipeTypeModel,
ResourceFile,
SinkGCSModel,
SinkKafkaModel,
SinkModel,
SinkS3Model,
Expand Down Expand Up @@ -600,7 +601,7 @@ def parse_pipe_file(resource: ResourceFile) -> PipeModel:
copy_mode = value
elif key == "EXPORT_SERVICE":
normalized = parse_quoted_value(value).lower()
if normalized not in {"kafka", "s3"}:
if normalized not in {"kafka", "s3", "gcs_hmac"}:
raise MigrationParseError(
resource.file_path,
"pipe",
Expand Down Expand Up @@ -721,7 +722,7 @@ def parse_pipe_file(resource: ResourceFile) -> PipeModel:
)

has_kafka_directives = export_topic is not None
has_s3_directives = any(
has_blob_directives = any(
value is not None
for value in (
export_bucket_uri,
Expand All @@ -731,16 +732,19 @@ def parse_pipe_file(resource: ResourceFile) -> PipeModel:
)
)

if has_kafka_directives and has_s3_directives:
if has_kafka_directives and has_blob_directives:
raise MigrationParseError(
resource.file_path,
"pipe",
resource.name,
"Sink pipe cannot mix Kafka and S3 export directives.",
)

# S3 and GCS sinks share the same EXPORT_* directive shape, so when EXPORT_SERVICE
# is omitted we default to "s3" (matching the SDK's historical behavior); GCS
# sinks must set EXPORT_SERVICE gcs_hmac explicitly to be recognized.
inferred_service = export_service or (
"kafka" if has_kafka_directives else "s3" if has_s3_directives else None
"kafka" if has_kafka_directives else "s3" if has_blob_directives else None
)
if not inferred_service:
raise MigrationParseError(
Expand All @@ -751,7 +755,7 @@ def parse_pipe_file(resource: ResourceFile) -> PipeModel:
)

if inferred_service == "kafka":
if has_s3_directives:
if has_blob_directives:
raise MigrationParseError(
resource.file_path,
"pipe",
Expand Down Expand Up @@ -793,7 +797,7 @@ def parse_pipe_file(resource: ResourceFile) -> PipeModel:
topic=export_topic,
schedule=export_schedule,
)
else:
elif inferred_service in {"s3", "gcs_hmac"}:
if has_kafka_directives:
raise MigrationParseError(
resource.file_path,
Expand All @@ -814,15 +818,34 @@ def parse_pipe_file(resource: ResourceFile) -> PipeModel:
"S3 sinks require EXPORT_BUCKET_URI, EXPORT_FILE_TEMPLATE, EXPORT_FORMAT, and EXPORT_SCHEDULE.",
)

sink = SinkS3Model(
service="s3",
connection_name=export_connection_name,
bucket_uri=export_bucket_uri,
file_template=export_file_template,
format=export_format,
schedule=export_schedule,
strategy=export_strategy, # type: ignore[arg-type]
compression=export_compression, # type: ignore[arg-type]
if inferred_service == "gcs_hmac":
sink = SinkGCSModel(
service="gcs_hmac",
connection_name=export_connection_name,
bucket_uri=export_bucket_uri,
file_template=export_file_template,
format=export_format,
schedule=export_schedule,
strategy=export_strategy, # type: ignore[arg-type]
compression=export_compression, # type: ignore[arg-type]
)
else:
sink = SinkS3Model(
service="s3",
connection_name=export_connection_name,
bucket_uri=export_bucket_uri,
file_template=export_file_template,
format=export_format,
schedule=export_schedule,
strategy=export_strategy, # type: ignore[arg-type]
compression=export_compression, # type: ignore[arg-type]
)
else:
raise MigrationParseError(
resource.file_path,
"pipe",
resource.name,
f'Unsupported EXPORT_SERVICE in strict mode: "{inferred_service}"',
)

params: list[PipeParamModel]
Expand Down
5 changes: 4 additions & 1 deletion src/tinybird_sdk/migrate/runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -254,7 +254,10 @@ def run_migrate(options: MigrateOptions | dict[str, Any]) -> MigrationResult:
)
continue

if pipe.sink and sink_connection_type != pipe.sink.service:
# The "gcs_hmac" sink service exports via a plain "gcs" connection.
sink_service = pipe.sink.service if pipe.sink else None
expected_connection_type = "gcs" if sink_service == "gcs_hmac" else sink_service
if pipe.sink and sink_connection_type != expected_connection_type:
errors.append(
MigrationError(
file_path=pipe.file_path,
Expand Down
14 changes: 13 additions & 1 deletion src/tinybird_sdk/migrate/types.py
Original file line number Diff line number Diff line change
Expand Up @@ -160,7 +160,19 @@ class SinkS3Model:
compression: Literal["none", "gzip", "snappy"] | None = None


SinkModel = SinkKafkaModel | SinkS3Model
@dataclass(frozen=True, slots=True)
class SinkGCSModel:
service: Literal["gcs_hmac"]
connection_name: str
bucket_uri: str
file_template: str
format: str
schedule: str
strategy: Literal["create_new", "replace"] | None = None
compression: Literal["none", "gzip", "snappy"] | None = None


SinkModel = SinkKafkaModel | SinkS3Model | SinkGCSModel


@dataclass(frozen=True, slots=True)
Expand Down
1 change: 1 addition & 0 deletions src/tinybird_sdk/schema/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,7 @@
CopyConfig,
KafkaSinkConfig,
S3SinkConfig,
GCSSinkConfig,
SinkConfig,
NodeDefinition,
)
Expand Down
27 changes: 24 additions & 3 deletions src/tinybird_sdk/schema/pipe.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@
from dataclasses import dataclass, field
from typing import Any

from .connection import KafkaConnectionDefinition, S3ConnectionDefinition
from .connection import GCSConnectionDefinition, KafkaConnectionDefinition, S3ConnectionDefinition
from .datasource import ColumnDefinition, DatasourceDefinition, SchemaDefinition, get_column_type
from .params import ParamValidator
from .token import TokenDefinition
Expand Down Expand Up @@ -70,7 +70,18 @@ class S3SinkConfig:
compression: SinkCompression | None = None


SinkConfig = KafkaSinkConfig | S3SinkConfig
@dataclass(frozen=True, slots=True)
class GCSSinkConfig:
connection: GCSConnectionDefinition
bucket_uri: str
file_template: str
format: str
schedule: str
strategy: SinkStrategy | None = None
compression: SinkCompression | None = None


SinkConfig = KafkaSinkConfig | S3SinkConfig | GCSSinkConfig


@dataclass(frozen=True, slots=True)
Expand Down Expand Up @@ -299,7 +310,17 @@ def _normalize_sink_config(raw: dict[str, Any]) -> SinkConfig:
strategy=raw.get("strategy"),
compression=raw.get("compression"),
)
raise ValueError("Sink connection must be a Kafka or S3 connection definition.")
if isinstance(connection, GCSConnectionDefinition):
return GCSSinkConfig(
connection=connection,
bucket_uri=raw["bucket_uri"],
file_template=raw["file_template"],
format=raw["format"],
schedule=raw["schedule"],
strategy=raw.get("strategy"),
compression=raw.get("compression"),
)
raise ValueError("Sink connection must be a Kafka, S3, or GCS connection definition.")


def define_sink_pipe(name: str, options: dict[str, Any]) -> PipeDefinition:
Expand Down
Loading
Loading