DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.
...
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
...