{"record":{"id":"eb1bbaa1642e07e7","repo":"apache/seatunnel","slug":"deserialize-message-failed-skip-this-message-mes","errorCode":null,"errorMessage":"Deserialize message failed, skip this message, message: {}","messagePattern":"Deserialize message failed, skip this message, message: (.+?)","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/source/KafkaRecordEmitter.java","lineNumber":98,"sourceCode":"                List<String> kafkaHeaderFields = consumerMetadata.getKafkaHeaderFields();\n                Collector<SeaTunnelRow> targetCollector =\n                        kafkaHeaderFields.isEmpty()\n                                ? outputCollector\n                                : headerInjectingCollector(\n                                        outputCollector,\n                                        consumerRecord.headers(),\n                                        kafkaHeaderFields);\n                ((CompatibleKafkaConnectDeserializationSchema) deserializationSchema)\n                        .deserialize(consumerRecord, targetCollector);\n            } else if (deserializationSchema instanceof NativeKafkaConnectDeserializationSchema) {\n                ((NativeKafkaConnectDeserializationSchema) deserializationSchema)\n                        .deserialize(consumerRecord, outputCollector);\n            } else {\n                deserializationSchema.deserialize(consumerRecord.value(), outputCollector);\n            }\n        } catch (Exception e) {\n            if (this.messageFormatErrorHandleWay == MessageFormatErrorHandleWay.SKIP) {\n                logger.warn(\n                        \"Deserialize message failed, skip this message, message: {}\",\n                        new String(consumerRecord.value()));\n            } else {\n                throw e;\n            }\n        }\n        // consumerRecord.offset + 1 is the offset commit to Kafka and also the start offset\n        // for the next run\n        splitState.setCurrentOffset(consumerRecord.offset() + 1);\n    }\n\n    private static Collector<SeaTunnelRow> headerInjectingCollector(\n            Collector<SeaTunnelRow> delegate, Headers headers, List<String> headerFieldNames) {\n        return new Collector<SeaTunnelRow>() {\n            @Override\n            public void collect(SeaTunnelRow record) {\n                delegate.collect(appendHeaderFields(record, headers, headerFieldNames));\n            }","sourceCodeStart":80,"sourceCodeEnd":116,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/source/KafkaRecordEmitter.java#L80-L116","documentation":"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.","triggerScenarios":"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.","commonSituations":"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.","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"],"exampleFix":"// before\nKafkaSource = {\n  message_format_error_handle_way = SKIP   # silent data loss\n}\n// after\nKafkaSource = {\n  message_format_error_handle_way = EXCEPTION  # fail fast on bad records\n}","handlingStrategy":"validation","validationCode":"// Pre-validate topic payload format with a consumer probe before running the job\nConsumerRecord<String, byte[]> r = pollOne(consumer);\ntry { schema.deserialize(r.value(), new TestCollector()); }\ncatch (Exception e) { throw new IllegalStateException(\"Topic payload incompatible with configured format\", e); }","typeGuard":null,"tryCatchPattern":"try {\n    deserializationSchema.deserialize(record.value(), collector);\n} catch (Exception e) {\n    // SKIP: record dropped (potential data loss) — count and alert\n    skippedMessages.inc();\n}","preventionTips":["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"],"tags":["kafka","deserialization","data-loss"],"backgroundTag":"deserialization-failed","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"}