---
title: "Kafka Streams, Connect and Flink"
description: "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."
---

> Queen MQ documentation, for AI agents
> Complete self-contained summary of Queen MQ: https://queenmq.com/llms-brief.txt
> Fetch that first when the question is about the product rather than about this page.
> Index of all pages: https://queenmq.com/llms.txt

# Kafka Streams, Connect and Flink

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:

```text
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](/reference/kafka/#topic-configs).

## 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.

## 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.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](/guides/kafka/), or see the [Kafka reference](/reference/kafka/) for the
topic configs and the APIs.

Source: https://queenmq.com/guides/kafka-frameworks/index.mdx
