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

  1. Verify the conflict_key option exactly matches a column name present in the row schema (case-sensitive).
  2. Check upstream transforms for schema changes that drop or rename the key column.
  3. Filter or coalesce null key rows before the sink (e.g. use a transform to replace null keys or drop invalid records).
  4. 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

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


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)