Skip to content

About

Ultra reliable & scalable ELT. Pull data from largest API sources with a fault-tolerant framework designed for billions records.

Resources

Contributing

Stars

2 stars

Watchers

1 watching

Forks

Repository files navigation

bizon ⚡️

Extract and load your largest data streams with a framework you can trust for billions of records.

Bizon is a lightweight Python ETL framework built around a checkpointed producer → queue → consumer pipeline. It extracts data from API clients, databases and event streams, optionally transforms records in flight with Python, and loads them into your warehouse — with native fault-tolerance, resumable checkpoints, and very high throughput, while staying small enough to read end to end.


Table of Contents


Features

  • Natively fault-tolerant — a checkpointing mechanism tracks progress so an incremental or stream pipeline resumes from its last committed cursor after a crash or restart. A full_refresh pipeline deliberately restarts instead, since it republishes the whole table.
  • High throughput — designed to process billions of records, using Polars DataFrames for memory-efficient, vectorized buffering and Parquet for batch loads.
  • Queue-system agnostic — a single Queue interface that any broker can be written against. The in-process Python queue is the maintained implementation; the RabbitMQ and Kafka/Redpanda adapters are currently unmaintained (see Queue types).
  • Pluggable connectors — 10 built-in sources and 5 destinations, all behind clean AbstractSource / AbstractDestination interfaces. Sources are auto-discovered — no registration needed.
  • Multiple sync modes — full_refresh, incremental (append-only), and continuous stream.
  • Secret resolution — keep secrets out of YAML with gsm:// (Google Secret Manager) and env:// references, resolved before validation so connectors only ever see plain strings.
  • In-pipeline transforms — apply user-defined Python to records in flight.
  • Pipeline metrics — exhaustive metrics with Datadog & OpenTelemetry tracing: ETAs, records processed, completion %, and source ↔ destination latency.
  • Lightweight & lean — a minimal codebase with few core dependencies (requests, pyyaml, pydantic, sqlalchemy, polars, pyarrow).

Architecture

Bizon uses a producer–consumer pattern with pluggable components:

YAML Config → RunnerFactory → Producer → Queue → Consumer → Destination
                                  ↑                    ↓
                                Source              Backend (checkpoints)
  • The Producer pulls records from the Source in iterations and pushes them onto the Queue as Polars DataFrames.
  • The Consumer pulls from the queue, applies Transforms, and writes batches to the Destination.
  • The Backend persists cursors so syncs are resumable.

Checkpointing & recovery. The producer writes its source cursor to the backend every syncCursorInDBEvery iterations, and the consumer records a destination cursor after each successful write. On restart, an incremental or stream pipeline reads the last destination cursor and resumes from the next iteration. Bizon's delivery contract is at-least-once — on recovery a batch may be re-written, so destinations are designed to tolerate duplicate writes.

A full_refresh run is not resumed. It republishes the whole table, so a partially-fetched run has nothing worth continuing: an interrupted job is retired and the next run starts a fresh one, staging into a clean temp table. (Resuming one would leave the published table stuck at its previous contents indefinitely, since the run never reaches the final iteration that swaps the table in.)

Abstraction Base Class Location
Source AbstractSource bizon/source/source.py
Destination AbstractDestination bizon/destination/destination.py
Queue AbstractQueue bizon/engine/queue/queue.py
Backend AbstractBackend bizon/engine/backend/backend.py
Runner AbstractRunner bizon/engine/runner/runner.py

Installation

Requires Python ≥ 3.10, < 3.13.

Note: package metadata currently declares >=3.9, but bizon does not import on 3.9 — parts of the codebase use X | None annotations that Python 3.9 evaluates at import time. Use 3.10 or newer.

pip install bizon

Optional features are installed via extras:

Extra pip install Enables
postgres bizon[postgres] PostgreSQL backend
bigquery bizon[bigquery] BigQuery backend & destinations (incl. Storage Write API)
kafka bizon[kafka] Kafka source (Avro/Schema Registry). Also the Kafka queue adapter, which is unmaintained
rabbitmq bizon[rabbitmq] RabbitMQ queue adapter (unmaintained)
gsheets bizon[gsheets] Google Sheets source
datadog bizon[datadog] Datadog metrics & tracing
secretmanager bizon[secretmanager] gsm:// Google Secret Manager references

Combine extras as needed, e.g. pip install 'bizon[bigquery,kafka,secretmanager]'.

For Development

# Install uv (if not already installed)
pip install uv

# Clone and install with all extras and dependency groups
git clone https://github.com/bizon-data/bizon-core.git
cd bizon-core
uv sync --all-extras --all-groups

# Run tests
uv run pytest tests/

# Format & lint (Ruff, line length 120)
uv run ruff format .
uv run ruff check --fix .

Quickstart

This example needs zero external dependencies — it uses the in-process Python queue and a SQLite backend (both defaults), reading from the built-in dummy source and writing to logger.

Create config.yml:

name: demo-creatures-pipeline

source:
  name: dummy
  stream: creatures
  authentication:
    type: api_key
    params:
      token: dummy_key

destination:
  name: logger
  config:
    dummy: dummy

Run it:

bizon run config.yml

CLI Reference

The CLI entry point is bizon (bizon.cli.main:cli).

Command Description
bizon run <config.yml> Run a pipeline from a YAML config
bizon source list List available sources and their streams
bizon stream list <source> List a source's streams, flagged [Supports incremental] / [Full refresh only]
bizon stream reset <config.yml> Queue a stream reset for the next run of that pipeline
bizon secrets check <config.yml> Dry-run every gsm:// / env:// reference and report (masked) results
bizon config validate <config.yml> Validate a config, including the source's own config class, without running it
bizon destination Subcommand group (no subcommands yet)

bizon run

bizon run config.yml \
  --custom-source ./my_source.py \   # Custom Python file implementing a Bizon source
  --runner thread \                  # thread | process | stream (default: thread)
  --log-level INFO \                 # DEBUG | INFO | WARNING | ERROR | CRITICAL
  --env-file .env \                  # Load env vars from a .env file (auto-detected if omitted)
  --result-json result.json          # Write the run's outcome as JSON

The --runner flag overrides engine.runner.type in the config; --log-level overrides engine.runner.log_level.

--result-json writes the run's outcome as JSON, so a wrapper can act on it without parsing logs:

{
  "status": "failure",
  "failure_class": "source",
  "job_id": "2d876897b416403f892e9f411b69efd6",
  "producer": "source_error",
  "consumer": "source_error",
  "stream": null,
  "records_written": 1200,
  "started_at": "2026-10-08T09:46:55.094362Z",
  "duration_s": 42.1,
  "error": "RuntimeError: page 2 came back empty"
}
  • status is success or failure.
  • failure_class is one of config, source, destination, backend, queue, transform, stream, killed or unknown.
  • error is the first error the run logged.
  • The file reads "status": "running" from the moment the run starts. A process killed from outside (OOM, a pod deadline) therefore leaves running behind, not the previous run's outcome.
  • The exit code is unchanged.

bizon secrets check

bizon secrets check config.yml [--env-file .env]

Resolves every reference in the config (without running the pipeline) and prints each one with a masked ✓/✗ status. Exits non-zero if any reference fails — handy as a pre-deploy gate.

bizon config validate

bizon config validate config.yml [--env-file .env] [--strict]

Validates the config the way a run would, without running anything:

  • the engine schema;
  • the stream name;
  • the source's own config class.

It also warns about two things that are not errors:

  • source: keys the source's config class does not declare, which a run silently ignores. This usually means a typo, or a key that belongs under destination.config.
  • deprecated keys.

References are resolved, so the command needs the same access as a run. In CI without secret access, pass --skip-references to validate them as literal strings instead. Use --strict to fail on warnings too.

Configuration Reference

A pipeline is defined by a single YAML file validated against BizonConfig (bizon/common/models.py). Unknown top-level keys are rejected.

Key Required Description
name ✅ Unique name identifying this sync
id — Stable identity the backend keys jobs, cursors and resets under. Defaults to name. Set it to the current name before renaming a pipeline, or the rename starts it from scratch
source ✅ Source connector config (see below)
destination ✅ Destination connector config; routed by its name field
transforms — List of in-pipeline transforms (default [])
engine — Backend, queue and runner config (sensible defaults if omitted)
secrets — Provider defaults for gsm:// / env:// resolution
monitoring — Datadog metrics & tracing
alerting — Slack alerting
streams — Multi-table routing (requires source.sync_mode: stream)

source keys

Common SourceConfig fields (bizon/source/config.py); each connector adds its own:

Key Default Description
name — Connector name (e.g. hubspot, notion, kafka)
stream — Stream to sync
sync_mode full_refresh full_refresh | incremental | stream
cursor_field None Timestamp field for incremental filtering (e.g. updated_at)
authentication None Auth block (type + params); connector-specific
force_ignore_checkpoint false Ignore existing checkpoints and restart from iteration 0. Redundant for full_refresh, which always starts fresh
reset false Re-fetch the whole stream and replace the destination table, then resume incremental (details)
max_iterations None Cap iterations per run (default: run until source is exhausted)
api_config.retry_limit 10 Retries before giving up on an API call
source_file_path None Path to a custom source file (same as --custom-source)
http None Opt-in HTTP policy for the default session (details)

HTTP retries and timeouts

Without an http block, the default session keeps its historical policy:

  • there is no timeout;
  • urllib3 only retries 413, 429 and 503, and only when the response carries a Retry-After header;
  • a 500, 502 or 504 fails the request immediately.

Setting the block, even as http: {}, switches the source to these defaults:

source:
  http:
    timeout: [10, 60]            # (connect, read) seconds, for requests that do not pass their own
    raise_for_status: true       # false for sources that inspect non-2xx responses themselves
    retries:
      total: 10
      backoff_factor: 2
      status_forcelist: [429, 500, 502, 503, 504, 520, 522, 524]
      allowed_methods: [GET, POST]
      retry_after_max: 120       # cap on a Retry-After wait; null honours the header as sent

Once retries run out, the request raises a requests.HTTPError carrying the last response, not a bare RetryError.

To tune retries in code, override get_retry_policy() and return a Retry (or bizon.source.session.CappedRetry for the Retry-After cap). Overriding get_session() replaces the session entirely: that drops the raise_for_status hook, and http is then ignored with a warning.

Annotated example

name: hubspot contacts to bigquery

source:
  name: hubspot
  stream: contacts
  properties:
    strategy: all          # connector-specific: fetch all properties
  authentication:
    type: api_key
    api_key: ${env://HUBSPOT_API_KEY}   # inline secret reference

destination:
  name: bigquery
  config:
    project_id: my-gcp-project-id
    dataset_id: bizon_test
    dataset_location: US
    create_dataset: false
    gcs_buffer_bucket: bizon-buffer
    gcs_buffer_format: parquet
    buffer_size: 10            # in MB; an iteration larger than this is written in buffer-sized chunks
    buffer_flush_timeout: 300  # in seconds

engine:
  backend:
    type: bigquery
    config:
      database: my-gcp-project
      schema: bizon_backend
      syncCursorInDBEvery: 10

Full, copy-pasteable examples live next to each connector under bizon/connectors/**/config/*.example.yml.

Connectors

Sources

All sources are auto-discovered. Run bizon source list to see what's installed and bizon stream list <source> for a source's streams.

Source Connects to Example streams Auth Incremental Extra
hubspot HubSpot CRM (v3) contacts, companies, deals api_key, oauth — —
notion Notion API pages, databases, data_sources, blocks, blocks_markdown, users (+ all_*) api_key ✅ —
kafka Kafka / Redpanda topic (Avro + Schema Registry) basic (SASL) stream kafka
gsheets Google Sheets worksheet service account, ADC — gsheets
cycle Cycle (GraphQL) customers api_key — —
periscope Periscope / Sisense charts, dashboards, views, users, databases cookies — —
sana_ai Sana AI Insight Reports insight_report oauth — —
pokeapi PokéAPI (public) pokemon, berry, item none — —
gbif GBIF biodiversity (public) occurrence none — —
dummy Mock source (testing/demos) creatures, plants api_key, oauth — —

notion is the reference incremental source — see bizon/connectors/sources/notion/src/source.py for get_records_after().

Destinations

Destination Writes to Sync modes Notable features
bigquery BigQuery full / incremental / stream Batch loads via GCS + Parquet; atomic table swaps via free copy jobs; async/batched load jobs; partitioning; optional schema unnesting
bigquery_streaming BigQuery full / incremental / stream Legacy streaming insert API; dynamic schema evolution; large-row fallback to load jobs
bigquery_streaming_v2 BigQuery full / incremental / stream Storage Write API (protobuf); higher throughput; multi-threaded appends; schema caching
file Local NDJSON file full / incremental / stream Atomic finalize on full refresh; append on incremental; optional unnesting
logger stdout (loguru) full / incremental / stream Logs records — for testing & debugging

Adding connectors. Sources are auto-discovered — drop a connector under bizon/connectors/sources/{name}/src/ and it appears. Destinations register in three places (see the guide). Use the /new-source and /new-destination Claude skills, or read docs/contributing/adding-sources.md and docs/contributing/adding-destinations.md.

Sync Modes

Mode Behavior
full_refresh Re-syncs all data from scratch on each run (default)
incremental Syncs only new/updated records since the last successful run (append-only)
stream Continuous streaming for real-time data (e.g. Kafka)

Incremental Sync

Incremental sync fetches only new or updated records since the last successful run, using an append-only strategy.

Configuration

source:
  name: your_source
  stream: your_stream
  sync_mode: incremental
  cursor_field: updated_at  # The timestamp field to filter records by

How It Works

┌─────────────────────────────────────────────────────────────────────┐
│                        INCREMENTAL SYNC FLOW                        │
├─────────────────────────────────────────────────────────────────────┤
│                                                                     │
│  1. Producer checks for last successful job                         │
│     └─> Backend.get_last_successful_stream_job()                    │
│                                                                     │
│  2. If found, creates SourceIncrementalState:                       │
│     └─> last_run = previous_job.created_at                          │
│     └─> cursor_field = config.cursor_field (e.g., "updated_at")     │
│                                                                     │
│  3. Calls source.get_records_after(source_state, pagination)        │
│     └─> Source filters: WHERE cursor_field > last_run               │
│                                                                     │
│  4. Records written to temp table: {table}_incremental              │
│                                                                     │
│  5. finalize() appends temp table to main table                     │
│     └─> BigQuery: free copy job (WRITE_APPEND), creating the         │
│         main table on first run; other destinations: INSERT INTO     │
│     └─> Deletes temp table                                          │
│                                                                     │
│  FIRST RUN: No previous job → falls back to get() (full refresh)    │
│                                                                     │
└─────────────────────────────────────────────────────────────────────┘

Supported Sources

Sources must implement get_records_after() to support incremental sync:

Source Cursor Field Notes
notion last_edited_time Supports pages, databases, blocks, blocks_markdown streams
(others) Varies Check source docs or implement get_records_after()

Supported Destinations

Destinations must implement finalize() with incremental logic:

Destination Support Notes
bigquery ✅ Append-only via temp table
bigquery_streaming_v2 ✅ Append-only via temp table
file ✅ Appends to existing file
logger ✅ Logs completion

Example: Notion Incremental Sync

name: notion pages incremental sync

source:
  name: notion
  stream: pages           # Options: databases, data_sources, pages, blocks, users
  sync_mode: incremental
  cursor_field: last_edited_time
  authentication:
    type: api_key
    params:
      token: secret_xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx
  database_ids:
    - "xxxxxxxx-xxxx-xxxx-xxxx-xxxxxxxxxxxx"
  page_size: 100

destination:
  name: bigquery
  config:
    project_id: my-gcp-project
    dataset_id: notion_data
    dataset_location: US
    gcs_buffer_bucket: my-gcs-bucket
    gcs_buffer_format: parquet

engine:
  backend:
    type: bigquery
    config:
      database: my-gcp-project
      schema: bizon_backend
      syncCursorInDBEvery: 2

First Run Behavior

On the first incremental run (no previous successful job):

  • Falls back to the get() method (full-refresh behavior)
  • All data is fetched and loaded
  • The job is marked successful
  • Subsequent runs use get_records_after() with the last_run timestamp

Source state and the run window

source_state carries more than last_run, which is timezone-aware UTC:

  • run_started_at is the current job's start time. Use it as the upper bound of your query: the next run's last_run is this exact value, so the windows tile with no gap or overlap.
  • state is whatever the previous successful run returned as SourceIteration.next_state. For example, a max cursor value or a per-partition offset; it must be JSON-serializable. The last non-None value of a run is persisted when the job succeeds. A run that emits none keeps the previous state, and a reset starts from {}.

Stream Reset

A reset re-fetches the whole stream once and replaces the destination table, then resumes incremental from that run. Use it when the table has drifted — a backfill, a bug in a transform, or records the source changed without bumping its cursor field.

There are three ways to ask for one; all do the same thing:

# One-shot, for a run you launch yourself
bizon run config.yml --reset

# Queued in the backend, for a pipeline whose command line is fixed by a scheduler.
# The next `bizon run config.yml` picks it up — no change to the cron/Airflow job.
bizon stream reset config.yml
bizon stream reset config.yml --cancel   # changed your mind

# One config templated across streams? Pick which stream to reset.
bizon stream reset config.yml --stream deals
source:
  sync_mode: incremental
  reset: true      # deprecated: resets on every run until removed

What the run does differently:

  • Producer ignores the watermark and calls get() instead of get_records_after()
  • Destination materializes the run as a full refresh, replacing the table rather than appending to it (for bigquery: staging into {table}_temp, then a WRITE_TRUNCATE copy job)
  • The job row stays incremental, so the next run uses this run as its new watermark

Notes:

  • Scoped to a single stream. The request is keyed on (name, source_name, stream_name) — the same triple as the watermark it overrides — so resetting one stream never affects another, even under the same pipeline name. bizon stream reset takes the stream from the config; --stream overrides it.
  • Only meaningful with sync_mode: incremental; ignored (with a warning) for the other modes.
  • Supported by every destination with a working full-refresh path, since that is what a reset materializes as. The exception is bigquery_streaming, which has no staging table and appends even on a full refresh; a reset there is rejected at config validation rather than duplicating your data.
  • A reset that crashes is retried as a reset — the request stays bound to its job, so a rerun cannot silently degrade into an append.
  • bizon stream reset writes to the backend, so it must reach the same one the pipeline uses. With the default file-based sqlite backend that means the same machine and file.

Engine Configuration

The engine block configures three pluggable subsystems. All have defaults, so engine is optional for local runs.

Backends (state storage)

The backend stores Bizon's state (jobs and cursors). Configured under engine.backend.

Type Use case
sqlite File-based SQLite — local development (default)
sqlite_in_memory Memory-only SQLite — unit tests
postgres PostgreSQL — production, frequent cursor updates (bizon[postgres])
bigquery BigQuery — lightweight production state storage (bizon[bigquery])

syncCursorInDBEvery (default 10) controls how often the source cursor is flushed to the backend — lower values mean finer-grained recovery, higher values mean less write overhead.

engine:
  backend:
    type: postgres
    config:
      host: localhost
      port: 5432
      database: bizon
      schema: bizon
      username: ${env://PG_USER}
      password: ${env://PG_PASSWORD}
      syncCursorInDBEvery: 10

On BigQuery the backend dataset must exist before the run starts — check_prerequisites() runs in init_job, before any destination code. The destination's create_dataset therefore cannot bootstrap it, even when both point at the same dataset. Set create_schema on the backend instead:

engine:
  backend:
    type: bigquery
    config:
      database: my-gcp-project
      schema: bizon_backend
      create_schema: true      # create the dataset if it is missing (default false)
      schema_location: EU      # permanent: BigQuery cannot move a dataset after creation

schema_location also pins the location of the backend's query jobs. Leave it unset to keep BigQuery's own default. When the backend and the destination share one dataset and both create flags are on, an unset schema_location is inherited from destination.dataset_location.

Queues

The queue carries records from the producer to the consumer. Configured under engine.queue.

Type Use case
python_queue In-process — the default, and the only maintained queue
rabbitmq RabbitMQ — unmaintained, does not currently run (bizon[rabbitmq])
kafka Kafka / Redpanda — unmaintained, does not currently run (bizon[kafka])

The producer pauses when the queue is full. By default "full" means about max_nb_messages records (1,000,000). To bound memory instead, set max_bytes: the producer then also pauses once the queued iterations take roughly that many bytes. The size is estimated from the average size of the iterations produced so far.

engine:
  queue:
    type: python_queue
    config:
      max_bytes: 500000000   # ~500 MB

The rabbitmq and kafka queue adapters are not supported and are known to be broken. Their consumers call write_records_and_update_cursor(source_records=...), but that method takes df_source_records, so they raise TypeError on the first message. They also predate the QUEUE_TERMINATION_ERROR signal added in 0.5.4, so they would not stop cleanly on a producer failure even once that is fixed. Use python_queue. The configuration below is retained for whoever revives them.

This is only about the queue adapters. The Kafka source (bizon/connectors/sources/kafka/) is maintained and covered by its own e2e workflow.

Spin up a broker locally with the provided compose files:

docker compose --file ./scripts/queues/kafka-compose.yml up      # Kafka
docker compose --file ./scripts/queues/redpanda-compose.yml up   # Redpanda
docker compose --file ./scripts/queues/rabbitmq-compose.yml up    # RabbitMQ
# Kafka
engine:
  queue:
    type: kafka
    config:
      queue:
        bootstrap_server: localhost:9092   # Kafka:9092 / Redpanda:19092

# RabbitMQ
engine:
  queue:
    type: rabbitmq
    config:
      queue:
        host: localhost
        queue_name: bizon

Runners

The runner controls how the producer and consumer execute. Configured under engine.runner.

Type Description Best for
thread ThreadPoolExecutor (default) I/O-bound sources (APIs, DBs)
process ProcessPoolExecutor, true parallelism CPU-bound sources/transforms
stream Single-threaded, inline Real-time stream sync & multi-stream routing

Transforms

Transforms apply user-defined Python to each record as it flows through the pipeline. Each transform receives the record's parsed payload as a data dict, mutates it, and the result is re-serialized. Transforms are applied sequentially, in order.

transforms:
  - label: debezium
    python: |
      from datetime import datetime

      cluster = data['value']['source']['name'].replace('_', '-')
      partition = data['partition']
      offset = data['offset']
      kafka_timestamp = datetime.utcfromtimestamp(
          data['value']['source']['ts_ms'] / 1000
      ).strftime('%Y-%m-%d %H:%M:%S.%f')

      deleted = False
      if data['value']['op'] == 'd':
          data = data['value']['before']
          deleted = True

A TransformModel has just two fields: label (a display name) and python (the code). There are no built-in transforms — the logic is entirely yours.

Secrets & References

Keep secrets out of YAML by referencing them with a URI scheme. Resolution runs once over the raw config before validation (bizon/engine/resolvers/), so connectors need no changes — they always read plain strings.

  • gsm://<id> → Google Secret Manager, latest version (ADC auth). Pin with gsm://<id>/versions/<N>, or pass a full gsm://projects/<p>/secrets/<id>/versions/<N> path. Requires bizon[secretmanager].
  • env://<VAR> → environment variable.
  • Inline form — embed in a larger string with ${...}, e.g. dsn: "postgres://u:${gsm://db-pw}@host/db" (multiple tokens allowed).
  • An optional secrets: block holds provider defaults (e.g. secrets.gsm.project_id).

Validate every reference before running:

bizon secrets check config.yml

Add a provider by dropping one adapter in bizon/engine/resolvers/adapters/ and one entry in _SCHEME_FACTORIES (bizon/engine/resolvers/resolver.py).

Monitoring & Alerting

Monitoring (bizon/monitoring/) — emit pipeline metrics and traces to Datadog. Install bizon[datadog].

monitoring:
  type: datadog
  config:
    enable_tracing: true
    datadog_agent_host: localhost   # or datadog_host_env_var: DD_AGENT_HOST
    datadog_agent_port: 8125
    tags:
      env: production
      team: data

Alerting (bizon/alerting/) — post alerts to Slack on configured log levels.

alerting:
  type: slack
  log_levels: [ERROR]              # defaults to [ERROR]
  config:
    webhook_url: ${env://SLACK_WEBHOOK_URL}

Multi-Stream Routing

For one-to-many pipelines (common with Kafka CDC), the streams block maps multiple source streams to their own destination tables and schemas in a single run. It requires source.sync_mode: stream and the stream runner.

source:
  name: kafka
  stream: topic
  sync_mode: stream
  # topics are auto-extracted from the streams block below
  bootstrap_servers: your-kafka-broker:9092
  group_id: your-consumer-group

destination:
  name: bigquery_streaming_v2
  config:
    project_id: your-gcp-project
    dataset_id: your_dataset
    dataset_location: US
    unnest: true

streams:
  - name: users
    source:
      topic: cdc.public.users
    destination:
      table_id: your-gcp-project.your_dataset.users
      clustering_keys: [id]
      record_schema:
        - { name: id, type: INTEGER, mode: REQUIRED }
        - { name: email, type: STRING, mode: NULLABLE }
  - name: orders
    source:
      topic: cdc.public.orders
    destination:
      table_id: your-gcp-project.your_dataset.orders
      clustering_keys: [id, user_id]

engine:
  runner:
    type: stream

See bizon/connectors/sources/kafka/config/kafka_streams.example.yml for a complete example.

BigQuery Partitioning

All three BigQuery destinations take the same time_partitioning block. Tables bizon creates are partitioned on _bizon_loaded_at by day unless you say otherwise.

destination:
  name: bigquery
  config:
    project_id: my-gcp-project
    dataset_id: bizon_test
    gcs_buffer_bucket: bizon-buffer
    time_partitioning:
      type: DAY              # DAY | HOUR | MONTH | YEAR
      field: _bizon_loaded_at
    enforce_partitioning: false
  • time_partitioning: DAY (a bare window) is still accepted and means the same as the block above.
  • field must exist in the destination schema. With unnest: true the table holds only the columns from record_schemas, so the _bizon_* metadata columns are not available — point field at one of your own TIMESTAMP/DATE/DATETIME columns. This is validated when the config loads.
  • field: null selects ingestion-time partitioning.

Repartitioning an existing table

BigQuery cannot change a table's partitioning after it is created. A copy job into an existing table keeps that table's spec, and CREATE OR REPLACE TABLE is rejected outright when the spec differs. So a table created by an older bizon (or by hand) stays as it is, however you configure time_partitioning. Bizon detects this and warns at publish time, naming both specs.

The only fix is to rebuild. Set enforce_partitioning: true and run a full refresh: the destination drops the table so it is recreated with the configured layout. Off by default because the table is briefly absent, and its description, labels, table ACLs and policy tags are not recreated.

For an incremental stream, a run stages only its own delta, so it never rebuilds even with the flag on — rebuilding from a delta would discard history. Use bizon stream reset with the flag enabled instead: the reset run re-fetches the stream, publishes as a full refresh (rebuilding the table), and leaves the job incremental so the next run resumes normally. This needs a source that can re-fetch the whole stream; otherwise repartition the table by hand.

Documentation & Contributing

Claude Code skills are available for common workflows: /new-source, /new-destination, and /run-checks.

License

Bizon is released under the GNU General Public License v3.0. See LICENSE.

About

Ultra reliable & scalable ELT. Pull data from largest API sources with a fault-tolerant framework designed for billions records.

Resources

Contributing

Stars

2 stars

Watchers

1 watching

Forks

Releases

Packages

Used by

Contributors

Languages