Skip to content

Kafka Streams, Connect and Flink

What Kafka Streams, Kafka Connect and Flink's Kafka connector ask of a broker, what Queen does with each request, what has run against it, and where each one stops.

Updated View as Markdown

Kafka Streams, Kafka Connect and Flink ask more of a broker than produce and consume. They create topics with configs and check them afterwards, keep their state in compacted topics, start reading from a point in time and commit transactions that outlive a process. Since 2.0.0-beta.5 the Kafka facade answers each of those, and Kafka Streams and Kafka Connect run against a three-node Queen cluster unchanged.

What has run

On the code of 2.0.0-beta.5, Kafka 4.3.1’s own distribution ran against a fresh three-node Queen cluster, through the facade in cluster mode, four times for each check:

  • Kafka Streams’ WordCountDemo: a group formed by the Streams assignor, a repartition topic and a compacted changelog, and every final count exact.
  • Kafka Connect in distributed mode: the worker created its config, offset and status topics (compacted, with 1, 25 and 5 partitions), and a file source fed a topic that a file sink wrote back out, every line accounted for.

Flink has not been run against Queen yet. What its Kafka connector asks for is covered below, one request at a time.

Topics created with configs

Streams creates its internal topics with segment.bytes, retention.ms=-1 and message.timestamp.type=CreateTime on repartition topics and cleanup.policy=compact on changelogs, and Connect creates its three topics compacted and then reads their configs back to check. A broker that refuses any of those keeps the framework from starting, so the facade accepts every topic config it can honour and says plainly what it does with each:

  • cleanup.policy, retention.ms, min.insync.replicas (up to the raft majority) and message.timestamp.type=CreateTime mean what they mean on Kafka, mapped onto the queue.
  • segment.bytes, segment.ms, retention.bytes, max.message.bytes and the compaction tuning keys are recorded and reported back as set, each with a line in DescribeConfigs’ documentation that says nothing enforces it. None of them can cost a reader a record.
  • Every other topic config is refused with INVALID_CONFIG, and so is a value Kafka itself would refuse, because accepting a setting that is not in force is a promise the broker cannot keep.

message.timestamp.type=LogAppendTime is one of the refusals, with a sentence that says why:

message.timestamp.type=LogAppendTime is not supported: a record fetched through this facade
carries the producer's timestamp, so `CreateTime` is the one timestamp type there is.
`LogAppendTime` would tell a consumer it reads the broker's clock

A topic created on one node with cleanup.policy=compact describes like this:

Config Value Source
cleanup.policy compact DYNAMIC_TOPIC_CONFIG
min.insync.replicas 1 DEFAULT_CONFIG
retention.ms -1 DEFAULT_CONFIG

The full list, with the values each key accepts, is in the reference.

Compacted topics are kept whole

Compaction guarantees a reader one thing: the last value of every key is in the log. Queen does not compact yet, so a topic with cleanup.policy=compact keeps every record. Retention is switched off on its queue, explicitly, and stays off whatever retention.ms says, so the last value of every key is always there, next to every value before it. The readers that depend on compacted topics replay them from the start and fold each record into state, which is what a Streams changelog restore and Connect’s config, offset and status topics do, and folding the whole log gives the same state as folding a compacted one, tombstones included, because each key’s records are in the log’s order.

The cost is size. The topic grows without bound, a restore reads all of history where Kafka would read one record per key, and a tombstone stays for good where Kafka removes it after delete.retention.ms. Real compaction as a policy any queue can have is designed in protocols/queen-kafka/COMPACTION.md and not built yet.

Starting from a point in time

offsetsForTimes, Flink’s OffsetsInitializer.timestamp(...) and every tool’s reset to a date ask for the first offset at or after a time. The facade answers from the broker’s POST /api/v1/fetch/offsets, a binary search over the times at which the partition’s records were appended, which with the default dedup index reads no payload. Produce three records a second apart, look up half a second after the first, and the answer is offset 1. A time later than the last append answers Kafka’s “no offset”: offsetsForTimes returns null for that partition, and Flink starts it at the end.

The clock is the broker’s, the moment each record was appended, and the record you then fetch carries the producer’s own timestamp. For an ordinary producer the two agree to within its linger and the network. They do not for a producer that sets timestamps itself, such as a replay of old events or a Streams application forwarding an input record’s time: then the answer is where the log was at that moment, a good place to start reading, but not the first record whose own timestamp is that time.

Kafka Streams

WordCountDemo, which keeps state and repartitions, runs with the default processing.guarantee=at_least_once. The Streams assignor needs nothing from the broker beyond the classic group protocol, so keep Streams on its default rebalance protocol. A few things to know before you run your own application:

  • Streams purges consumed records from its repartition topics with DeleteRecords, which Queen does not offer yet, so those topics keep everything (Streams creates them with retention.ms=-1). You can give them a retention with kafka-configs.sh --alter, as long as it comfortably outlasts the longest your application can fall behind.
  • exactly_once_v2 runs on Kafka transactions, which need a facade outside cluster mode, and it has not been run against Queen yet.
  • When you try it against a fresh cluster, use a new application.id or an empty state.dir. Streams restores local state from its state directory, and state left over from an earlier cluster is counted again: in WordCountDemo, every count doubles.

Kafka Connect

A distributed worker needs nothing Queen-specific in its configuration. It creates its compacted config, offset and status topics, checks that they are compacted, and replays them when it restarts, so connector configurations and source offsets survive a worker restart as they do on Kafka. Connect asks for replication factor 3 by default, which is accepted and reported as 1, because replication is raft’s: every voter holds every partition. Source and sink connectors are ordinary producers and consumers. Exactly-once source support (exactly.once.source.support=enabled) runs on transactions, so it needs a facade outside cluster mode, and it has not been run against Queen yet.

KafkaSource is a plain consumer: it reads partitions, commits offsets, and finds its start offsets with ListOffsets, including the timestamp lookup above. KafkaSink with AT_LEAST_ONCE is an ordinary producer. With EXACTLY_ONCE the sink commits a Kafka transaction per checkpoint, and three details decide whether that works:

  • Flink asks for a transaction.timeout.ms of one hour, which a stock Kafka broker refuses until its transaction.max.timeout.ms is raised. The facade’s ceiling is one hour by default (QUEEN_KAFKA_TXN_MAX_TIMEOUT_MS), so Flink’s default is accepted as it is.
  • After a task manager fails over, Flink resumes the checkpoint’s transaction from a new connection. The facade keeps a transaction’s stage past the connection that opened it, until its own timeout, so that commit lands. It does not survive a restart of the node holding the stage: the resumed commit is refused with INVALID_TXN_STATE, and Flink drops that checkpoint’s records.
  • Cluster mode refuses transactions, so an exactly-once sink needs a facade outside it.

Limits

Besides the ones above, two group features are not supported: static membership (group.instance.id) and the KIP-848 consumer group protocol. All three frameworks leave both off unless you turn them on, so leave them at their defaults.

Back to Kafka clients, or see the Kafka reference for the topic configs and the APIs.

Navigation

Type to search…

↑↓ navigate↵ selectEsc close