Bug Description
Severity: CRITICAL
Component: PreparedStatementExecutor, ClickHouseBatchRunnable
Version: 2.9.1
Problem
The connector has no transaction atomicity guarantee. When a batch of records is processed, it is split into partitions via Lists.partition(), and each partition is executed independently via ps.executeBatch(). If a later partition fails after earlier partitions have already committed, there is no rollback mechanism for the already-committed partitions.
Additionally, MySQL transaction boundaries (BEGIN/COMMIT/ROLLBACK) from Debezium CDC events are not tracked or preserved. Records from different MySQL transactions can be mixed in a single batch, breaking transactional consistency.
Affected Code
File: sink-connector/src/main/java/com/altinity/clickhouse/sink/connector/db/batch/PreparedStatementExecutor.java
// Lines 155-261: executePreparedStatement()
Lists.partition(entry.getValue(), (int)maxRecordsInBatch).forEach(batch -> {
try (PreparedStatement ps = metadata.getPreparedStatement(conn, insertQuery)) {
for (ClickHouseStruct record : batch) {
// ... add records to batch ...
ps.addBatch();
}
int[] batchResult = ps.executeBatch(); // Line 236: Commits this partition
result.set(true);
} catch (Exception e) {
// Line 251: Previous partitions already committed!
// No rollback possible for those partitions
throw new RuntimeException(e);
}
});
File: sink-connector/src/main/java/com/altinity/clickhouse/sink/connector/executor/ClickHouseBatchRunnable.java
// Line 634: queryToRecordsMap is a plain HashMap, no transaction context
Map<MutablePair<String, Map<String, Integer>>,
List<ClickHouseStruct>> queryToRecordsMap = new HashMap<>();
Impact
- Data inconsistency: Partial batches committed to ClickHouse while remaining records are lost
- No rollback: ClickHouse does not support traditional ROLLBACK for MergeTree tables
- Cross-transaction mixing: Records from different MySQL transactions may be mixed in a single ClickHouse batch
- Duplicate records on retry: When a partial batch fails and retries, already-committed records are re-inserted
Reproduction
-- MySQL transaction
BEGIN;
INSERT INTO orders VALUES (1, 'pending');
INSERT INTO orders VALUES (2, 'pending');
UPDATE orders SET status = 'confirmed' WHERE id = 1;
COMMIT;
-- If connector crashes mid-batch:
-- ClickHouse may have INSERT(1) and INSERT(2) but NOT the UPDATE
-- Result: order 1 shows 'pending' instead of 'confirmed'
Proposed Fix
- Track transaction boundaries: Parse Debezium transaction metadata (
__debezium.transactions topic) to group records by MySQL transaction
- All-or-nothing batch execution: Execute all partitions of a batch within a single ClickHouse transaction context (for supported engines)
- Idempotent inserts: Use ReplacingMergeTree with version columns to handle retry scenarios
- Offset management: Only commit Kafka offsets after ALL partitions of a batch succeed
Bug Description
Severity: CRITICAL
Component: PreparedStatementExecutor, ClickHouseBatchRunnable
Version: 2.9.1
Problem
The connector has no transaction atomicity guarantee. When a batch of records is processed, it is split into partitions via
Lists.partition(), and each partition is executed independently viaps.executeBatch(). If a later partition fails after earlier partitions have already committed, there is no rollback mechanism for the already-committed partitions.Additionally, MySQL transaction boundaries (BEGIN/COMMIT/ROLLBACK) from Debezium CDC events are not tracked or preserved. Records from different MySQL transactions can be mixed in a single batch, breaking transactional consistency.
Affected Code
File:
sink-connector/src/main/java/com/altinity/clickhouse/sink/connector/db/batch/PreparedStatementExecutor.javaFile:
sink-connector/src/main/java/com/altinity/clickhouse/sink/connector/executor/ClickHouseBatchRunnable.javaImpact
Reproduction
Proposed Fix
__debezium.transactionstopic) to group records by MySQL transaction