Skip to content

Streaming (dw-streaming)

The Streaming agent handles the operational core of Kafka-based pipelines: getting a stream topology configured correctly, and knowing when it’s falling behind. Describe a source and a sink — say, a Debezium Postgres source into a Snowflake sink — and it validates the topology and generates the Kafka Connect JSON configuration, converter classes included, with your partition count, replication factor, and SLA targets baked in.

Once a stream runs, the agent watches it: consumer lag per partition, connector status, latency, and error rates, aggregated into per-connector and overall stream health. When something drifts — lag climbing, throughput sagging against an SLA — it generates tuning recommendations. It recommends; you apply. Nothing is auto-tuned behind your back.

  • Kafka Connect config generation. configure_stream validates a stream topology and produces working Kafka Connect JSON for source and sink connectors — the real config generation logic, including converter classes.
  • SLA-aware setup. Pass SLA targets (max latency, max lag in records) at configuration time and they become the baseline the tuning engine judges against.
  • Consumer-lag monitoring. monitor_lag returns lag per partition for a topic, plus total lag across partitions — the first number you want when a downstream table goes stale.
  • Stream health aggregation. get_stream_health rolls connector status, latency, and error rates into per-connector and overall health for a topology.
  • Tuning recommendations. get_recommendations analyzes lag, throughput, and health against your SLAs and proposes concrete tuning changes — human-applied, never auto-applied.

“Generate a Kafka Connect config for a Debezium Postgres source into a Snowflake sink on topic orders-cdc, 12 partitions.”

“What’s the consumer lag on orders-cdc, per partition?”

“Give me the health of the orders-cdc topology — connectors, latency, error rates.”

“Lag has been climbing all afternoon. What tuning changes do you recommend?”

  • Kafka Schema Registry — the streaming entry in the connector catalog; schema registration and compatibility checks flow through it.
  • Warehouse sinks — configs commonly target sinks like Snowflake; connect the warehouse to verify the receiving end.

The agent starts in 🟡 Evaluation, and config generation is real logic even there: configure_stream produces valid Kafka Connect JSON for your topology before any credential exists. Lag and health monitoring run against built-in sample data until your systems are connected and verified. See Verify your setup.

  • This is a deliberately small surface — configuration, lag, health, tuning. It manages Kafka Connect topologies; it is not a stream processor and doesn’t run your Flink or Spark jobs.
  • Recommendations are never auto-applied. If you want a hands-off tuner, this agent will disappoint you on purpose: it proposes, you decide.
  • In 🟡 Evaluation, lag and health numbers describe the sample topology, not your cluster. Generated configs are real, but validate them against your environment before deploying.
  • Monitoring depth depends on what’s connected; where a metric isn’t available from your setup, the agent says so rather than inventing a number.