apache/seatunnel · error · KafkaConnectorException
COMMON_ILLEGAL_ARGUMENT (CommonErrorCodeDeprecated.ILLEGAL_ARGUMENT)
COMMON_ILLEGAL_ARGUMENT (CommonErrorCodeDeprecated.ILLEGAL_ARGUMENT)
Error message
Field name { %s } is not found! What it means
Kafka sink's DefaultSeaTunnelRowSerializer.topicExtractor supports a '{field_name}' topic pattern. When the referenced field does not exist in the SeaTunnelRowType of the incoming rows, it throws KafkaConnectorException with code ILLEGAL_ARGUMENT ('Field name { %s } is not found!'). Note the message interpolates the whole topic string, not the extracted field name.
Solutions
- Add the referenced field to the upstream data so it exists in the row schema.
- Correct the field name in the topic config to match an actual field in the row's SeaTunnelRowType (exact, case-sensitive match).
- Use a static topic string if you do not need per-row topic routing.
Example fix
# before
topic = "events-{topic_name}" # no field 'topic_name' in schema
# after
topic = "events-{topic}" # 'topic' is an actual row field Defensive patterns
Strategy: validation
Validate before calling
// validate at job-config time
String topicPattern = "events-{partition}";
java.util.regex.Matcher m = java.util.regex.Pattern.compile("\\{(\\w+)\\}").matcher(topicPattern);
if (m.find()) {
String field = m.group(1);
if (!java.util.Arrays.asList(rowType.getFieldNames()).contains(field)) {
throw new IllegalArgumentException("topic field not in schema: " + field);
}
} Try / catch
try {
serializer = DefaultSeaTunnelRowSerializer.create(...);
} catch (KafkaConnectorException e) {
if (e.getSeaTunnelErrorCode().getCode().equals(CommonErrorCodeDeprecated.ILLEGAL_ARGUMENT.getCode())) {
// fix topic pattern or add missing field, then rebuild serializer
}
throw e;
} Prevention
- Verify every {field} in the topic pattern exists in the source schema (case-sensitive).
- Update topic patterns whenever upstream column names change.
- Prefer static topics when per-row routing is not required.
When it happens
Trigger: Configuring sink topic as a pattern like 'logs-{partition}' where 'partition' is not a field of the SeaTunnel row schema (rowType.getFieldNames() does not contain it).
Common situations: Typo in the field name inside the topic pattern; renamed upstream columns without updating the Kafka sink topic config; using a pattern field that exists only in a different table's schema in multi-table jobs.
Understand the failure class
Background: "is required", "must be set", "missing required field": configuration validation errors across open-source libraries — this error's family across 36 libraries.
Related errors
- COMMON_ILLEGAL_ARGUMENT (CommonErrorCodeDeprecated.ILLEGAL_ARGUMENT)
- COMMON-02
- COMMON_ERROR_CODE-1
- COMMON_UNSUPPORTED_DATA_TYPE (CommonErrorCodeDeprecated.UNSUPPORTED_DATA_TYPE)
- COMMON_UNSUPPORTED_DATA_TYPE (CommonErrorCodeDeprecated.UNSUPPORTED_DATA_TYPE)
AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10).
Data as JSON: /api/errors/fd494c61a1574fd7.
Report an issue: GitHub.
Appendix: source
Thrown at seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/serialize/DefaultSeaTunnelRowSerializer.java:258
|| MessageFormat.NATIVE.equals(format))
&& topic == null) {
int topicFieldIndex =
rowType.indexOf(CompatibleDebeziumJsonDeserializationSchema.FIELD_TOPIC);
return row -> row.getField(topicFieldIndex).toString();
}
String regex = "\\$\\{(.*?)\\}";
Pattern pattern = Pattern.compile(regex, Pattern.DOTALL);
Matcher matcher = pattern.matcher(topic);
boolean isExtractTopic = matcher.find();
if (!isExtractTopic) {
return row -> topic;
}
String topicField = matcher.group(1);
List<String> fieldNames = Arrays.asList(rowType.getFieldNames());
if (!fieldNames.contains(topicField)) {
throw new KafkaConnectorException(
CommonErrorCodeDeprecated.ILLEGAL_ARGUMENT,
String.format("Field name { %s } is not found!", topic));
}
int topicFieldIndex = rowType.indexOf(topicField);
return row -> {
Object topicFieldValue = row.getField(topicFieldIndex);
if (topicFieldValue == null) {
throw new KafkaConnectorException(
CommonErrorCodeDeprecated.ILLEGAL_ARGUMENT, "The column value is empty!");
}
return topicFieldValue.toString();
};
}
private static Function<SeaTunnelRow, byte[]> keyExtractor(
List<String> keyFields,
SeaTunnelRowType rowType,
MessageFormat format,View on GitHub (pinned to cf67b549a7)