Kafka ksqlDB Query Builder

Build a ksqlDB statement with the rules attached. Push against pull, co-partitioning, the mandatory WITHIN on a stream-stream join, and the KEY_FORMAT default that catches everybody, each stated next to the choice that triggers it.

What to build
Input
Join

Both sides must be co-partitioned: same key AND same partition count. ksqlDB does not check the partition count and produces missing rows rather than an error.

Aggregation

One row per window when it closes. Windowed aggregations only, and it delays every result by the grace period.

Output

Ad hoc queries only. Without it the query is a pull query, which can only read a materialised TABLE by key.

query.sql

updates as you type

    Common mistakes

    These are the ones that fail silently. The config is accepted, nothing raises an error, and the consequence arrives later.

    1. Creating a STREAM over a topic that should be a TABLE

      A stream is an append-only log of events and a table is the latest value per key. Choosing wrongly gives either duplicate rows or lost history, and the fix is a rewrite.

      Instead:Events are a stream. State keyed by an identifier is a table.

    2. Joining a stream to a stream with no window

      A stream-stream join requires a WITHIN clause, because without one the join has unbounded state.

      Instead:Specify the window, and size it against the real skew between the two streams.

    3. Expecting a persistent query to survive a schema change

      A CREATE ... AS SELECT query is bound to the schema at creation. An incompatible change to the source breaks it, and it must be dropped and recreated.

      Instead:Plan schema evolution around the queries, and use compatible changes where possible.

    It looks like SQL, and the rules that matter have no SQL equivalent

    Nearly every ksqlDB mistake is a rule a database does not have. Six of them decide whether a statement runs at all, and none of them is visible in the syntax.

    Push and pull are the same keyword away from each other

    EMIT CHANGES makes a query a push query: it never completes, it streams results for as long as the client is connected. Without it the query is a pull query, which returns rows and finishes. A pull query can only read a materialised TABLE, and only by key: a stream has no state to read, so there is no way to make one answer a pull query. A CREATE ... AS SELECT is always a push query, which is why the statement below always carries EMIT CHANGES.

    A join needs the same key and the same partition count

    Co-partitioning is two conditions and ksqlDB only checks one of them. The keys have to match, which it verifies. The partition counts have to match too, which it does not: a join between a six-partition stream and a twelve-partition table produces missing rows and no error at all. That is the hardest ksqlDB bug to find, because the query runs and the output looks plausible.

    A stream-stream join must be windowed

    There is no unbounded stream-stream join, because holding every record from both sides forever is not a thing a system can do. WITHIN is mandatory and it also decides how much state each instance keeps, so a wide window is a memory and disk cost paid continuously rather than a setting.

    PARTITION BY is a repartition, which is the most expensive thing here

    Changing the key means the data has to be redistributed, so ksqlDB writes it to a new topic and reads it back. That is a full round trip through the broker per record: serialise, produce, replicate, fetch, deserialise. It is often necessary and it is never cheap, and it is worth knowing you asked for it.

    KEY_FORMAT is not VALUE_FORMAT

    They are separate settings and the key defaults to KAFKA, which means a raw serialised primitive with no schema at all. Somebody who set VALUE_FORMAT='AVRO' and expects the key to be Avro too gets a key that a schema-aware deserialiser cannot read. The statement below always writes KEY_FORMAT explicitly, so the choice is visible rather than inherited.

    What this cannot see

    Your schemas, so it cannot tell you a column does not exist. Your partition counts, so it can state the co-partitioning rule and not whether you satisfy it. Whether the output topic already exists, which matters because PARTITIONS and REPLICAS are ignored silently when it does. It checks the statement's structure, which is what can be checked without a server.

    More kafka tools

    Kafka Confluent Wire Format Decoder The five junk bytes in front of your payload Kafka Key to Partition Mapper Which partition does this key land on? Kafka Topic Name Validator Legal, risky, or 249 characters too long? Kafka Replication Safety Checker How many brokers can you lose Kafka Producer Config Linter Will it start, and will it lose a record? Kafka Message Payload Decoder The first five bytes are usually not data Kafka Connect Source Connector Generator tasks.max is a ceiling, not a count Kafka Connect Sink Connector Generator A dead letter queue with no context headers is a pile of records Kafka Connect SMT Chain Builder The order is the transforms list Kafka MirrorMaker 2 Config Generator It renames every topic by default Kafka Partition Reassignment Generator The throttle is not optional Strimzi Kafka Resource Generator Without the cluster label, nothing happens Kafka mTLS Config Generator The certificate is the identity Kafka Schema Registry Config Generator The compatibility direction is your deployment order Kafka Exactly-Once Config Generator Half of it is worse than none Kafka Broker and KRaft Config Generator The internal topics that break a one-broker cluster Kafka Quota Generator Byte rates are per broker, not per cluster Kafka Streams Config Generator application.id is four things at once Kafka Connect Worker Config Generator Security three times, or the tasks fail Kafka Retention and Unit Converter log.retention.hours does not take milliseconds Kafka Timestamp Converter Two sentinels and two meanings Kafka .properties to YAML Converter Dotted keys stay flat Kafka Streams Internal Topic Predictor Create them before Streams does Kafka ACL Generator The grant you forgot is on another resource type Kafka Topic Config Generator min.insync.replicas is the one that matters Kafka client.properties Generator The file every CLI tool asks for Kafka Producer Config Generator No password field, on purpose Kafka Consumer Config Generator The commit mode decides the semantics Kafka Disk and Retention Calculator retention.bytes is per partition Kafka Partition Count Calculator The number you can never reduce Kafka Cluster Sizing Calculator The traffic no client metric shows Kafka Consumer Lag Catch-Up Calculator Whether it ever clears, not just when Kafka Producer Batching Calculator linger.ms=0 still batches Kafka Segment and Index Sizing Why retention.ms is a lower bound Kafka Rebalance Duration Estimator What a rolling restart really costs Kafka Cost Estimator Your rates, so nothing goes stale Kafka Config Explorer by Version The answer depends on the release Kafka Default Config Reference What moved under a config you never edited Kafka OAuth Bearer Token Decoder Will Kafka accept it, and can it refresh Kafka Record Header Viewer Headers are a list, not a map Kafka Topic Regex Subscription Tester Kafka matches the whole name Kafka ACL Permission Matrix Viewer DENY beats every ALLOW Kafka Connect Config Validator The mistakes that raise no error Kafka Consumer Group Id Validator Which broker coordinates the group Kafka Partition Assignment Visualizer Leadership is the load, not replicas Kafka Consumer Assignment Visualizer The three assignors disagree Kafka ZooKeeper to KRaft Config Converter The authorizer class nobody changes Kafka Config to Strimzi Half of it belongs elsewhere Kafka Docker Compose Generator (KRaft) Reachable from inside and outside Kafka JAAS Config Decoder The line that stops SASL working Kafka CRC32C Calculator Which CRC, over which bytes Kafka Config Upgrade Checker What breaks when you upgrade Kafka Kafka Config Diff Which change actually changed something Kafka Consumer Config Linter Why the group rebalances, and where the records went Kafka Avro Schema Validator The defaults Avro accepts and rejects Kafka Schema Compatibility Checker What the registry will say, before you ask it Kafka Avro Schema Diff Which direction each change breaks Kafka Compression Comparison Measured on your bytes Kafka Delivery Semantics Exactly-once has a consumer half Kafka Connect SMT Predicate Tester negate reads backwards Kafka Streams Topology Viewer Count the repartitions Kafka Connect Pipeline Visualizer The order things really run in Kafka Protobuf Binary Decoder Works without the .proto Kafka Protobuf JSON Converter Why your JSON does not round-trip Kafka Protobuf to Avro Schema What does not survive the conversion Kafka Avro Binary Decoder Wrong schema, no error Kafka Avro JSON Converter Why the console producer rejects your line Kafka Avro Sample Data Generator Records that actually serialize Kafka JSON to Avro Schema What JSON cannot tell you Kafka JSON Schema to Avro What does not survive the conversion Kafka SASL JAAS Generator One login module, four syntaxes Kafka CLI Command Builder kcat is librdkafka, not Kafka

    Elsewhere on the site