Versions Compared

Key

  • This line was added.
  • This line was removed.
  • Formatting was changed.

...

No existing features are deprecated. This is purely additive.

5. Test Plan

5.1 Unit Tests

1. `WorkerShareSinkTaskTest`: Tests the core poll-put-acknowledge loop using a `MockShareConsumer`.
   - Verify ACCEPT after successful `task.put()`
   - Verify RELEASE after `RetriableException`
   - Verify REJECT after unrecoverable exception
   - Verify `commitSync()` is called at configured intervals

...

3. `SinkConnectorConfigTest` (modified): Validate the new configuration properties and their defaults.

5.2 Integration Tests

1. Basic Share Group Sink: Deploy a sink connector with `group.protocol=share` and verify all records are delivered.
2. Elastic Scaling: Start with 2 tasks, scale to 6, verify no records are lost and throughput increases.
3. Task Failure and Re-delivery: Kill a task mid-processing, verify records are re-delivered to surviving tasks within `acquisition.lock.timeout.ms`.
4. No Duplicate Loss: Produce N records, consume with at-least-once Share Group sink, verify received count >= N.
5. Interoperability: Verify that standard consumer group connectors and Share Group connectors can coexist in the same Connect cluster.

5.3 System Tests

1. Long-running throughput test: Measure throughput and latency of Share Group vs. consumer group sink connectors under sustained load.
2. Chaos test: Randomly kill tasks and brokers, verify zero data loss with at-least-once semantics.

6.

...

Rejected Alternatives

Alternative 1: Modify the SinkTask Interface to Add acknowledge()

We considered adding `acknowledge(SinkRecord)` and `release(SinkRecord)` methods to the `SinkTask` interface, giving connectors explicit control over acknowledgments. This was rejected because:
- It would break backward compatibility with all existing sink connectors
- Most connectors don't need per-record acknowledgment control
- The worker runtime can make correct acknowledgment decisions based on `put()` success/failure

Alternative 2: Use Share Groups Only for MirrorMaker2

We considered limiting Share Group support to the `MirrorSourceConnector` only, as Kafka-to-Kafka is the most obvious use case. This was rejected because:
- It would require changes to the MM2 `consumer.assign()` model, which is complex
- Generic sink connectors (e.g., JDBC, Elasticsearch, S3) benefit equally from elastic scaling
- Building it into the Connect runtime benefits all connectors automatically

Alternative 3: Exactly-Once from Day One

We considered requiring exactly-once semantics for the initial implementation. This was rejected because:
- KIP-1289 (transactional share acknowledgments) is not yet implemented
- At-least-once is sufficient for the majority of sink connector use cases
- Idempotent sinks (upsert to database, overwrite to S3) achieve effective exactly-once with at-least-once delivery
- Exactly-once can be added as a follow-up without breaking changes

...