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

  1. Set message_format_error_handle_way = EXCEPTION temporarily to surface the full stack trace
  2. Verify the configured deserialization format matches what producers write to the topic
  3. Dump the failing record (the message is printed in the log) and check schema compatibility
  4. 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

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


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)