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) andmessage.timestamp.type=CreateTimemean what they mean on Kafka, mapped onto the queue.segment.bytes,segment.ms,retention.bytes,max.message.bytesand 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 clockA 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 withkafka-configs.sh --alter, as long as it comfortably outlasts the longest your application can fall behind. exactly_once_v2runs 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.idor an emptystate.dir. Streams restores local state from its state directory, and state left over from an earlier cluster is counted again: inWordCountDemo, 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.
Flink
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.msof one hour, which a stock Kafka broker refuses until itstransaction.max.timeout.msis 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.