apache/seatunnel · error · IllegalArgumentException
Conflict key field ' ' value is null or not found in row
Error message
Conflict key field '${conflictKey}' value is null or not found in row What it means
When the sink operates in conflict-key (upsert/CDC) mode, getConflictKeyValue extracts the conflict key column's value from a SeaTunnelRow to build the conflict clause. If the row has no field matching the conflict key, or its value is null, an IllegalArgumentException is thrown because a dedup/upsert write cannot be identified.
Solutions
- Verify the conflict_key option exactly matches a column name present in the row schema (case-sensitive).
- Check upstream transforms for schema changes that drop or rename the key column.
- Filter or coalesce null key rows before the sink (e.g. use a transform to replace null keys or drop invalid records).
- Align source catalog schema with the Databend table so field names resolve correctly.
Example fix
// before
sink {
Databend {
conflict_key = "id_"
}
}
// after: use the actual key column present in the schema
sink {
Databend {
conflict_key = "id"
}
} Defensive patterns
Strategy: validation
Validate before calling
// before configuring the sink, ensure conflict_key is in the schema
Set<String> fields = Set.of(rowType.getFieldNames());
if (!fields.contains(conflictKey)) {
throw new IllegalArgumentException("conflict_key not in schema: " + conflictKey);
}
// and reject null keys upstream
if (row.getField(rowType.getFieldIndex(conflictKey)) == null) {
throw new IllegalArgumentException("Null conflict key in row: " + row);
} Type guard
boolean hasNonNullConflictKey(SeaTunnelRow row, SeaTunnelRowType type, String key) {
int i = type.getFieldIndex(key);
return i >= 0 && row.getField(i) != null;
} Try / catch
try {
writer.write(row);
} catch (IllegalArgumentException e) {
if (e.getMessage().startsWith("Conflict key field")) {
// route row to a dead-letter log instead of failing the job
} else throw e;
} Prevention
- Keep conflict_key exactly matching a schema column name
- Avoid upstream transforms that drop or rename the key column
- Filter or coalesce null keys before the sink
- Re-check conflict_key after any schema evolution change
When it happens
Trigger: conflict_key option names a column that does not exist in the row's field names, or the value at that index is null in the arriving row; row field order shifted after schema evolution so the lookup misses the key.
Common situations: Typo in the conflict_key config option; CDC rows where the key column was dropped by an upstream transform; schema mismatch between source and Databend table causing key column not found; null keys produced by a null-valued upstream column.
Understand the failure class
Background: "must not be empty", "cannot be empty" — required-field validation errors across open-source libraries — this error's family across 41 libraries.
Related errors
- Column already exists in table . Skipping add column…
- Column already exists in table . Skipping change column…
- Column does not exist in table . Skipping drop column…
- COMMON-06
- COMMON_UNSUPPORTED_OPERATION
AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10).
Data as JSON: /api/errors/05fa3b830cbcebef.
Report an issue: GitHub.
Appendix: source
Thrown at seatunnel-connectors-v2/connector-databend/src/main/java/org/apache/seatunnel/connectors/seatunnel/databend/sink/DatabendSinkWriter.java:478
/**
* Get the value of the conflict key field from the row. This value will be used as the ID in
* the raw table.
*/
private String getConflictKeyValue(SeaTunnelRow row) {
String[] fieldNames = catalogTable.getSeaTunnelRowType().getFieldNames();
int index = Arrays.asList(fieldNames).indexOf(conflictKey);
if (index >= 0 && index < row.getFields().length) {
Object value = row.getField(index);
if (value != null) {
return value.toString();
}
}
// This should not happen in a proper CDC setup where conflict key values are always present
// If we reach here, it indicates a data issue
throw new IllegalArgumentException(
"Conflict key field '" + conflictKey + "' value is null or not found in row");
}
private final ObjectMapper objectMapper = new ObjectMapper();
private String convertRowToJson(SeaTunnelRow row) {
try {
ObjectNode jsonNode = objectMapper.createObjectNode();
String[] fieldNames = catalogTable.getSeaTunnelRowType().getFieldNames();
Object[] fields = row.getFields();
for (int i = 0; i < fieldNames.length; i++) {
String fieldName = fieldNames[i];
Object value = fields[i];
if (value == null) {
jsonNode.putNull(fieldName);
} else if (value instanceof String) {View on GitHub (pinned to cf67b549a7)