Kafka Connect Pipeline Visualizer

Paste a connector configuration, as .properties or as the JSON body you would POST, and see the pipeline in the order it actually runs. Transforms execute in the order of the transforms list, not the order of the blocks in the file.

Paste below, or drop a file anywhere on this panel

Or drop a file anywhere on this panel. Nothing is uploaded: the analysis runs in this tab.

The answer appears here

Paste on the left and press Draw it. Nothing leaves this tab.

Examples

Real input you can load into the tool above. Each one shows a different thing going wrong, because that is what the tool is for.

Debezium with an SMT

The connector, its transform chain and where records go

{"name":"pg-source","config":{"connector.class":"io.debezium.connector.postgresql.PostgresConnector","tasks.max":"1","transforms":"unwrap","transforms.unwrap.type":"io.debezium.transforms.ExtractNewRecordState"}}

A sink with a dead letter queue

Where failed records go, and the settings that decide whether they are kept at all

{"name":"s3-sink","config":{"connector.class":"io.confluent.connect.s3.S3SinkConnector","tasks.max":"4","topics":"orders","errors.tolerance":"all","errors.deadletterqueue.topic.name":"dlq"}}

Common mistakes

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

  1. Setting tasks.max higher than the source can split

    A JDBC source with one table produces one task no matter what tasks.max says. The extra capacity is not used and the number is misleading.

    Instead:Match tasks.max to the real parallelism: partitions for a sink, tables or partitions for a source.

  2. Chaining transforms without checking order

    SMTs apply in the order named in the transforms list, and an ExtractNewRecordState after a field-level transform operates on a different record shape than expected.

    Instead:Name them in the order they should run, and test with a real record.

  3. Running a converter that does not match the data

    A JsonConverter reading Avro produces a deserialization error on every record, and the connector fails fast rather than skipping.

    Instead:Match key and value converters to what is actually on the topic, which are separate settings.

A source and a sink run their stages in opposite orders

This is the fact the picture exists for. A source runs its transforms and then the converter; a sink runs the converter and then its transforms. An SMT chain does not transplant between the two unchanged, and nothing in the config says so.

Which side of serialisation a transform sees

On a source connector the transform runs on the connector's own record, before the converter turns it into bytes, so it sees the connector's schema and types. On a sink connector the converter runs first, so the transform sees a record that has already been deserialised from the wire. A transform that reaches into a field by name works in both cases; one that depends on the type of that field often does not, because the converter is where the type is decided.

The chain runs in the order of the transforms list

transforms is a comma-separated list of aliases, and that list is the execution order. The transforms.<alias>.* blocks below it can appear in any order at all, and frequently do, because people group them by what they configure rather than by when they run. A chain that reads correctly top to bottom can run in a completely different order, and the only place the real order appears is that one line.

transforms=unwrap,mask,route
transforms.mask.type=...MaskField$Value
transforms.route.type=...RegexRouter
transforms.unwrap.type=...ExtractField$Value

A predicate gates a transform, it does not filter records

predicates defines them and transforms.<alias>.predicate attaches one to a single transform. When it does not match, that transform is skipped and the record carries on through the rest of the chain untouched. With negate=true the gate is inverted. A predicate defined and never attached does nothing whatsoever and raises no error, which is why this page reports one.

The dead letter queue needs three things at once

It only exists for sink connectors. It only receives anything when errors.tolerance is all. And it only catches converter and transform failures, never a failure inside the connector itself. Setting the topic and leaving the tolerance at its default none is the commonest way to end up with an empty DLQ and a stopped task, because the first bad record fails the task instead of being routed.

errors.tolerance=all
errors.deadletterqueue.topic.name=dlq
errors.deadletterqueue.context.headers.enable=true

What this cannot see

Whether the connector class exists or what its own settings mean, because those are the connector's business and there are hundreds of them. Whether a transform's configuration is valid for its type, for the same reason. The worker's converter defaults, so a config that sets no converter is drawn as inheriting one. It reads the direction from the class name, which is a convention rather than a rule, and says so when the name does not follow it.

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 ksqlDB Query Builder It looks like SQL and the rules are not Kafka Connect SMT Predicate Tester negate reads backwards Kafka Streams Topology Viewer Count the repartitions 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