{"record":{"id":"05fa3b830cbcebef","repo":"apache/seatunnel","slug":"conflict-key-field-conflictkey-value-is-null","errorCode":null,"errorMessage":"Conflict key field '${conflictKey}' value is null or not found in row","messagePattern":"Conflict key field '(.+?)' value is null or not found in row","errorType":"exception","errorClass":"IllegalArgumentException","httpStatus":null,"severity":"error","filePath":"seatunnel-connectors-v2/connector-databend/src/main/java/org/apache/seatunnel/connectors/seatunnel/databend/sink/DatabendSinkWriter.java","lineNumber":478,"sourceCode":"\n    /**\n     * Get the value of the conflict key field from the row. This value will be used as the ID in\n     * the raw table.\n     */\n    private String getConflictKeyValue(SeaTunnelRow row) {\n        String[] fieldNames = catalogTable.getSeaTunnelRowType().getFieldNames();\n        int index = Arrays.asList(fieldNames).indexOf(conflictKey);\n\n        if (index >= 0 && index < row.getFields().length) {\n            Object value = row.getField(index);\n            if (value != null) {\n                return value.toString();\n            }\n        }\n\n        // This should not happen in a proper CDC setup where conflict key values are always present\n        // If we reach here, it indicates a data issue\n        throw new IllegalArgumentException(\n                \"Conflict key field '\" + conflictKey + \"' value is null or not found in row\");\n    }\n\n    private final ObjectMapper objectMapper = new ObjectMapper();\n\n    private String convertRowToJson(SeaTunnelRow row) {\n        try {\n            ObjectNode jsonNode = objectMapper.createObjectNode();\n            String[] fieldNames = catalogTable.getSeaTunnelRowType().getFieldNames();\n            Object[] fields = row.getFields();\n\n            for (int i = 0; i < fieldNames.length; i++) {\n                String fieldName = fieldNames[i];\n                Object value = fields[i];\n\n                if (value == null) {\n                    jsonNode.putNull(fieldName);\n                } else if (value instanceof String) {","sourceCodeStart":460,"sourceCodeEnd":496,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-connectors-v2/connector-databend/src/main/java/org/apache/seatunnel/connectors/seatunnel/databend/sink/DatabendSinkWriter.java#L460-L496","documentation":"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.","triggerScenarios":"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.","commonSituations":"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.","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."],"exampleFix":"// before\nsink {\n  Databend {\n    conflict_key = \"id_\"\n  }\n}\n// after: use the actual key column present in the schema\nsink {\n  Databend {\n    conflict_key = \"id\"\n  }\n}","handlingStrategy":"validation","validationCode":"// before configuring the sink, ensure conflict_key is in the schema\nSet<String> fields = Set.of(rowType.getFieldNames());\nif (!fields.contains(conflictKey)) {\n    throw new IllegalArgumentException(\"conflict_key not in schema: \" + conflictKey);\n}\n// and reject null keys upstream\nif (row.getField(rowType.getFieldIndex(conflictKey)) == null) {\n    throw new IllegalArgumentException(\"Null conflict key in row: \" + row);\n}","typeGuard":"boolean hasNonNullConflictKey(SeaTunnelRow row, SeaTunnelRowType type, String key) {\n    int i = type.getFieldIndex(key);\n    return i >= 0 && row.getField(i) != null;\n}","tryCatchPattern":"try {\n    writer.write(row);\n} catch (IllegalArgumentException e) {\n    if (e.getMessage().startsWith(\"Conflict key field\")) {\n        // route row to a dead-letter log instead of failing the job\n    } else throw e;\n}","preventionTips":["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"],"tags":["databend","conflict-key","upsert","null-value"],"backgroundTag":"empty-required-field","analyzedSha":"cf67b549a7a6c35fa0beb12d83c62892427ea919","analyzedAt":"2026-09-10T21:44:55.265Z","contentChangedAt":"2026-09-10T21:44:55.265Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}