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
- 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.
- Verify the connector's schema (table_path and columns) exactly matches `DESCRIBE TABLE` output.
- Reproduce the SQL from the message in clickhouse-client to confirm the error is server-side, then adjust filters or table.
- 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
- Keep the connector schema in sync with the real table; re-run schema inference after DDL changes.
- Test pushed-down filters manually in clickhouse-client first.
- Check replica health for Distributed tables before large reads.
- Set realistic query timeouts for the data volume being read.
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
- Failed EXECUTE SQL in catalog
- Failed to analyze SQL complexity using EXPLAIN, fallback to…
- Failed to execute query
- Failed to get clickhouse server config
- Failed to parse SQL statement
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)