onkar2405 opened a new pull request, #23091:
URL: https://github.com/apache/kafka/pull/23091
## Problem
MirrorMaker 2 masks two failure modes that matter when a Kafka topic is used
as a write-ahead log
and replicated to a DR cluster.
**1. Silent data loss.** If the source topic's retention policy deletes
records before MM2 replicates
them, the consumer's next fetch resumes at a valid-but-later offset. The gap
in the replicated stream
is never reported.
**2. Source topic reset.** If a source topic is deleted and recreated, the
tracked offset becomes
invalid. MM2 rewinds to the earliest offset and re-replicates the new topic
on top of the previously
mirrored data.
Both share a root cause. MM2 applies `consumer.auto.offset.reset=earliest`
to the replication
consumer (`MirrorConnectorConfig#sourceConsumerConfig`), so the broker's
`OFFSET_OUT_OF_RANGE` error
is resolved inside the consumer and never reaches the connector. Even if it
did, the only place that
could see it — the broad `catch (KafkaException e)` in
`MirrorSourceTask#poll` — logs a warning and
returns `null`.
This PR adds opt-in offset validation: MM2 detects the invalid offset, logs
the topic, partition and
problematic offset, and fails the task with a purpose-built exception.
---
## Design Rationale
### Files and classes modified
**New**
| File | Purpose |
| --- | --- |
| `connect/mirror/.../mirror/DataLossException.java` | Unreplicated records
removed by retention |
| `connect/mirror/.../mirror/TopicResetException.java` | Source topic
deleted and recreated |
| `connect/mirror/.../mirror/MirrorSourceTaskOffsetValidationTest.java` |
Unit tests for detection and classification |
|
`connect/mirror/.../mirror/integration/MirrorConnectorsIntegrationOffsetValidationTest.java`
| End-to-end tests on embedded clusters |
**Modified**
| File | Change |
| --- | --- |
| `MirrorSourceConfig.java` | New `offset.validation.enabled` config;
`sourceConsumerConfig` override that sets `auto.offset.reset=none` |
| `MirrorSourceTask.java` | Flag plumbing; `seekToBeginning` for uncommitted
partitions; `OffsetOutOfRangeException` handling and classification |
| `MirrorSourceConfigTest.java` | Four tests covering the new config and its
effect on the consumer config |
Both exceptions extend `ConnectException`. Roughly 120 lines of production
code and 330 lines of
tests.
### Why this implementation strategy
**Let the consumer detect the problem, rather than checking offsets
ourselves.**
The alternative is to track the last replicated offset per partition and
compare it against the
offset of the next record returned by `poll()`, flagging a gap. That
duplicates bookkeeping the
consumer already does, only notices a problem once records start flowing
again — a purged *idle*
partition would never be flagged — and adds per-record cost on the hot path.
Setting `auto.offset.reset=none` makes the broker's own out-of-range signal
the trigger. It costs
nothing when replication is healthy and it fires on the first fetch after
the offset goes bad.
**Distinguish the two conditions by log start offset.**
Both surface as the same `OffsetOutOfRangeException`. The log start offset
of the affected partition
separates them cleanly:
| Log start offset | Meaning | Exception |
| --- | --- | --- |
| `> 0` | Records ahead of our position were deleted by retention |
`DataLossException` |
| `== 0` | The log begins at zero again — the topic was recreated |
`TopicResetException` |
`MirrorSourceTask#classifyOffsetOutOfRange` queries
`Consumer#beginningOffsets` for the affected
partitions and applies that rule. When a single batch contains both
conditions, data loss wins,
because it is the one with unrecoverable consequences on the target cluster.
If the `beginningOffsets`
lookup itself fails, the task still fails — as `TopicResetException` —
rather than falling back to the
silent path.
**Fail fast rather than auto-recover.**
Neither situation has a safe automatic remedy. Skipping the gap loses
records with no record of what
was lost; rewinding to `earliest` after a topic recreation duplicates or
interleaves data on the DR
cluster. Both are operator decisions. Failing the task surfaces the problem
through Connect's status
API with a message naming the topic, partition, requested offset and log
start offset; the operator
resumes by resetting the connector offsets via the Connect offsets API.
> Note: the task description asks for a `TopicResetException` and fail-fast
behaviour (Task 2), while
> the evaluation criteria mention "auto-recovery for topic resets". This PR
implements fail-fast per
> the explicit task requirement. Auto-recovery would be a small follow-up
given the detection logic is
> already in place.
**Ship it disabled by default.**
Enabling this unconditionally would change the failure semantics of every
existing MM2 deployment: a
cluster quietly tolerating retention-driven gaps would start failing tasks
after an upgrade.
`offset.validation.enabled` defaults to `false`, so upgrading is a no-op and
operators of WAL-style
topics opt in deliberately.
**Let validation override an explicit reset policy.**
If a user sets both `offset.validation.enabled=true` and
`consumer.auto.offset.reset=latest`, the
override wins and `none` is applied, because any other policy silently
defeats the feature. Rejecting
the combination with a `ConfigException` would also be defensible; this is
asserted in
`testOffsetValidationTakesPrecedenceOverExplicitAutoOffsetReset` so the
choice is explicit rather
than incidental.
### How it integrates with the MirrorMaker 2 lifecycle
| Hook | Change |
| --- | --- |
| `MirrorSourceConfig#sourceConsumerConfig` (override) | Sets
`auto.offset.reset=none` when enabled. The shared helper in
`MirrorConnectorConfig` is untouched, so `MirrorCheckpointConnector` and
`MirrorHeartbeatConnector` keep their existing consumer settings. |
| `MirrorSourceTask#start` | Reads the flag from the task config. |
| `MirrorSourceTask#initializeConsumer` | Collects partitions with no
committed offset and calls `seekToBeginning` on them. |
| `MirrorSourceTask#poll` | New `catch (OffsetOutOfRangeException)` placed
ahead of the existing `catch (KafkaException)`, so the specific case is handled
before the generic warn-and-continue path. |
**Preserving first-start behaviour.** With `auto.offset.reset=none` the
consumer has no fallback
position, so a partition MM2 has never replicated would fail on its first
poll. `initializeConsumer`
therefore seeks those partitions to the beginning explicitly. This
reproduces exactly what
`auto.offset.reset=earliest` did, but only for genuinely new partitions — an
offset that exists and
is invalid is no longer silently repaired.
No changes to connector startup, task assignment, offset syncs, checkpoints,
or heartbeats. Because
both exceptions extend `ConnectException`, the Connect runtime's existing
task-failure handling
applies unchanged; no new lifecycle hooks were required.
---
## Configuration
| Name | Type | Default | Importance | Description |
| --- | --- | --- | --- | --- |
| `offset.validation.enabled` | boolean | `false` | medium | Fail the task
with `DataLossException` / `TopicResetException` when the offset to replicate
from is no longer available on the source cluster. Sets
`consumer.auto.offset.reset=none` on the replication consumer. |
---
## Testing
**Unit — `MirrorSourceTaskOffsetValidationTest`** (mocked consumer)
- data loss when the log start offset is ahead of the requested offset
- topic reset when the log start offset is zero
- data loss takes precedence in a mixed batch
- a failed `beginningOffsets` lookup still fails the task
- with the flag off, `poll` returns `null` and never queries
`beginningOffsets` — default behaviour preserved
- `seekToBeginning` is called for uncommitted partitions only when the flag
is on
- deterministic partition ordering in the error message
**Unit — `MirrorSourceConfigTest`** (4 added)
- validation off by default, consumer still gets `auto.offset.reset=earliest`
- validation on yields `auto.offset.reset=none`
- validation overrides an explicitly configured reset policy
- other consumer configs (`max.poll.records`, `enable.auto.commit`) are
unaffected
**Integration — `MirrorConnectorsIntegrationOffsetValidationTest`**
(embedded Connect clusters)
- *Data loss*: replicate a topic, stop the source connector, rewind its
offsets to the start of each
partition, then `deleteRecords` past that point on the source cluster.
Restarting the connector
fails the task with `DataLossException`.
- *Topic reset*: replicate a topic, stop the source connector, delete and
recreate the source topic
with the same partition count, produce a short batch. Restarting the
connector fails the task with
`TopicResetException`.
Both drive the failure with `deleteRecords` and topic recreation rather than
waiting on retention, so
they are deterministic and not timing-dependent.
---
## AI Usage Documentation
AI was used substantially on this change, in an agentic setup (Claude, via
Anthropic's Cowork
desktop tool). Being specific about the split:
**What the AI did**
- Read the current `MirrorSourceTask`, `MirrorSourceConfig`,
`MirrorConnectorConfig`,
`MirrorSourceTaskConfig`, `MirrorSourceTaskTest` and
`MirrorConnectorsIntegrationBaseTest` on
`trunk`, and identified the two places responsible for the masking
behaviour — the
`putIfAbsent(AUTO_OFFSET_RESET_CONFIG, "earliest")` in
`sourceConsumerConfig` and the
`catch (KafkaException)` in `poll()`.
- Drafted all of the production code: the two exception classes,
`classifyOffsetOutOfRange`, the
`seekToBeginning` handling in `initializeConsumer`, and the config
plumbing.
- Drafted all of the tests, including the `deleteRecords` / topic-recreation
approach used to make
the integration scenarios deterministic.
- Verified every external API it relied on against this checkout rather than
from memory —
`OffsetOutOfRangeException#offsetOutOfRangePartitions`,
`RecordsToDelete#beforeOffset`,
`EmbeddedKafkaCluster#createTopic`/`createAdminClient`/`consume`,
`EmbeddedConnect#connectorStatus`/`stopConnector`/`resumeConnector`,
`ConnectorStateInfo.TaskState#trace`, and the `TestUtils#waitForCondition`
overloads.
- Checked whether the new config key would break existing assertions in
`MirrorSourceConfigTest`,
`MirrorMakerConfigTest` and `MirrorConnectorConfigTest`, and confirmed it
would not — which also
surfaced that `MirrorSourceConfig#sourceConsumerConfig` had no test
coverage, prompting the four
added config tests.
- Wrote this PR description.
**What the AI did not do**
- It never compiled or ran anything. Its sandbox had only JDK 11 and no
usable access to the repo
working tree, so every Gradle run — unit tests, integration tests and
checkstyle — was executed
locally by me.
- The design decisions were reviewed and accepted by me rather than taken on
trust: opt-in default,
fail-fast over auto-recovery, log start offset as the classification
signal, overriding the shared
`MirrorConnectorConfig` helper on the subclass instead of editing it, and
validation taking
precedence over an explicit `auto.offset.reset`. The AI flagged the last
one as a genuine
either-way choice rather than deciding it silently.
- The full diff was read and reviewed by me before commit.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]