apache/seatunnel · error · ClickhouseConnectorException

QUERY_DATA_ERROR

QUERY_DATA_ERROR

Error message

Query data with sql error. sql: %s, message: %s

What it means

ClickhouseProxy.batchFetchRecords executes a data query (used by rows()) and wraps ClickHouseException from the SELECT into QUERY_DATA_ERROR, including the SQL and the ClickHouse message. It indicates the data-reading query itself failed — bad SQL, type conversion problems, or a server error mid-read.

Solutions

  1. Read the embedded ClickHouse message in the error (the SQL and server reason are printed) and fix the query cause — usually a column/schema mismatch.
  2. Verify the connector's schema (table_path and columns) exactly matches `DESCRIBE TABLE` output.
  3. Reproduce the SQL from the message in clickhouse-client to confirm the error is server-side, then adjust filters or table.
  4. Increase query timeouts / check replica health if the failure is a timeout or dropped connection during large reads.

Example fix

// before
source {
  Clickhouse {
    table_path = "db.events"
    # schema lists column "user_id" but table has "uid"
  }
}
// after
source {
  Clickhouse {
    table_path = "db.events"
    schema = { columns = { uid = "string" } }
  }
}
Defensive patterns

Strategy: validation

Validate before calling

// validate configured columns against the live table before the read
String desc = "DESCRIBE TABLE " + tablePath.getFullName();
Set<String> actual = runQuery(desc).stream().map(r -> r[0]).collect(toSet());
configuredColumns.forEach(c -> {
    if (!actual.contains(c)) throw new IllegalArgumentException("column mismatch: " + c);
});

Try / catch

try {
    rows = proxy.rows(tablePath, rowType, ...);
} catch (ClickhouseConnectorException e) {
    LOG.error("data query failed: {}", e.getMessage()); // message already contains SQL + server reason
    throw e;
}

Prevention

When it happens

Trigger: Calling batchFetchRecords / rows() when the generated SELECT fails: schema mismatch between configured SeaTunnelRowType and actual table columns, invalid WHERE/filter expression, query timeout, or ClickHouse node failure during fetch.

Common situations: Configured columns don't match the real ClickHouse table (renamed/dropped columns); unsupported filter pushdown producing invalid SQL; querying a Distributed table whose underlying shard is down; version-specific clickhouse-client incompatibilities.

Understand the failure class

Background: "query failed", "%w: SQL error" — wrapped database query errors in Go libraries explained — this error's family across 3 libraries.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/e794735deb22ee9e. Report an issue: GitHub.

Appendix: source

Thrown at seatunnel-connectors-v2/connector-clickhouse/src/main/java/org/apache/seatunnel/connectors/seatunnel/clickhouse/util/ClickhouseProxy.java:493

        }
    }

    public List<SeaTunnelRow> batchFetchRecords(
            String sql, TablePath tablePath, SeaTunnelRowType seaTunnelRowType) {
        List<SeaTunnelRow> seaTunnelRowList = new ArrayList<>();
        log.debug("run query data sql: {}", sql);

        try (ClickHouseResponse response = clickhouseRequest.query(sql).executeAndWait()) {
            response.stream()
                    .forEach(
                            record -> {
                                SeaTunnelRow seaTunnelRow =
                                        ClickhouseUtil.convertToSeaTunnelRow(
                                                record, seaTunnelRowType, tablePath.getFullName());
                                seaTunnelRowList.add(seaTunnelRow);
                            });
        } catch (ClickHouseException e) {
            throw new ClickhouseConnectorException(
                    ClickhouseConnectorErrorCode.QUERY_DATA_ERROR,
                    String.format(
                            "Query data with sql error. sql: %s, message: %s", sql, e.getMessage()),
                    e);
        }

        return seaTunnelRowList;
    }

    public boolean isComplexSql(String sql) {
        try {
            String explainSql = "EXPLAIN " + sql;

            try (ClickHouseResponse response =
                    getClickhouseConnection().query(explainSql).executeAndWait()) {
                List<String> explainOutput =
                        response.stream()
                                .map(record -> record.getValue(0).asString())

View on GitHub (pinned to cf67b549a7)