Sitelet https://github.com/Altinity/clickhouse-sink-connector/issues/1261
Skip to content

[BUG-TX-1] No Transaction Atomicity - Partial Batch Commits on Failure #1261

Description

@minguyen9988

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

  1. Track transaction boundaries: Parse Debezium transaction metadata (__debezium.transactions topic) to group records by MySQL transaction
  2. All-or-nothing batch execution: Execute all partitions of a batch within a single ClickHouse transaction context (for supported engines)
  3. Idempotent inserts: Use ReplacingMergeTree with version columns to handle retry scenarios
  4. Offset management: Only commit Kafka offsets after ALL partitions of a batch succeed

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions