apache/seatunnel · warning
Deserialize message failed, skip this message, message
Error message
Deserialize message failed, skip this message, message: {} What it means
KafkaRecordEmitter failed to deserialize a consumer record. When the configured message_format_error_handle_way is SKIP, the record is dropped with a warning containing the raw message bytes instead of failing the task; otherwise the exception is rethrown and fails the reader.
Solutions
- Set message_format_error_handle_way = EXCEPTION temporarily to surface the full stack trace
- Verify the configured deserialization format matches what producers write to the topic
- Dump the failing record (the message is printed in the log) and check schema compatibility
- Align producer and consumer schemas, or upgrade the connector's format dependencies
Example fix
// before
KafkaSource = {
message_format_error_handle_way = SKIP # silent data loss
}
// after
KafkaSource = {
message_format_error_handle_way = EXCEPTION # fail fast on bad records
} Defensive patterns
Strategy: validation
Validate before calling
// Pre-validate topic payload format with a consumer probe before running the job
ConsumerRecord<String, byte[]> r = pollOne(consumer);
try { schema.deserialize(r.value(), new TestCollector()); }
catch (Exception e) { throw new IllegalStateException("Topic payload incompatible with configured format", e); } Try / catch
try {
deserializationSchema.deserialize(record.value(), collector);
} catch (Exception e) {
// SKIP: record dropped (potential data loss) — count and alert
skippedMessages.inc();
} Prevention
- Avoid message_format_error_handle_way=SKIP in production; monitor skipped-record metrics
- Enforce schema compatibility checks on the producer side
- Keep consumer format dependencies in sync with producer schema versions
When it happens
Trigger: emitRecord calls deserializationSchema.deserialize(...) and it throws (malformed payload, schema drift, incompatible format); config has message_format_error_handle_way = SKIP, so the catch block logs and continues to the next record.
Common situations: Producer upgraded to a newer schema the consumer doesn't understand; wrong value_format/deserializer configured; binary or non-JSON/Avro bytes on the topic; corrupted or truncated messages; missing Confluent schema registry dependencies.
Related errors
- COMMON-02
- COMMON-02
- CONVERT_TO_CONNECTOR_TYPE_ERROR_SIMPLE
- Please invoke DeserializationSchema#deserialize
- Please invoke DeserializationSchema#deserialize
AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10).
Data as JSON: /api/errors/eb1bbaa1642e07e7.
Report an issue: GitHub.
Appendix: source
Thrown at seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/source/KafkaRecordEmitter.java:98
List<String> kafkaHeaderFields = consumerMetadata.getKafkaHeaderFields();
Collector<SeaTunnelRow> targetCollector =
kafkaHeaderFields.isEmpty()
? outputCollector
: headerInjectingCollector(
outputCollector,
consumerRecord.headers(),
kafkaHeaderFields);
((CompatibleKafkaConnectDeserializationSchema) deserializationSchema)
.deserialize(consumerRecord, targetCollector);
} else if (deserializationSchema instanceof NativeKafkaConnectDeserializationSchema) {
((NativeKafkaConnectDeserializationSchema) deserializationSchema)
.deserialize(consumerRecord, outputCollector);
} else {
deserializationSchema.deserialize(consumerRecord.value(), outputCollector);
}
} catch (Exception e) {
if (this.messageFormatErrorHandleWay == MessageFormatErrorHandleWay.SKIP) {
logger.warn(
"Deserialize message failed, skip this message, message: {}",
new String(consumerRecord.value()));
} else {
throw e;
}
}
// consumerRecord.offset + 1 is the offset commit to Kafka and also the start offset
// for the next run
splitState.setCurrentOffset(consumerRecord.offset() + 1);
}
private static Collector<SeaTunnelRow> headerInjectingCollector(
Collector<SeaTunnelRow> delegate, Headers headers, List<String> headerFieldNames) {
return new Collector<SeaTunnelRow>() {
@Override
public void collect(SeaTunnelRow record) {
delegate.collect(appendHeaderFields(record, headers, headerFieldNames));
}View on GitHub (pinned to cf67b549a7)