A sink connector throws DataException / SerializationException on deserialize. How do you diagnose and resolve converter and schema mismatches?
answer
- mismatch: data format != converter config
- inspect raw bytes: 0x00+ID=registry, {schema,payload}=JSON envelope
- check key and value separately
- errors.tolerance=all + DLQ topic + context.headers
- verify schema.registry.url / schemas.enable
basics
~20 sThe sink's converter usually doesn't match how the data was actually written. Identify the real on-topic format, align the sink's key/value converter and settings (e.g. schemas.enable, schema.registry.url) to it, and use error tolerance / a DLQ for poison records.
solid answer
~50 sMost sink deserialization failures are a converter mismatch: the data was produced in one format but the sink reads it with another. Diagnose by inspecting the actual bytes — e.g. a leading 0x00 magic byte plus 4-byte ID means a registry-backed (Avro/Protobuf) message, plain { means JSON, etc. Check that key.converter and value.converter (and their schemas.enable / schema.registry.url settings) match how the producer wrote the data; remember key and value can differ. Classic cases: AvroConverter sink against plain-JSON data (fails on the magic byte); JsonConverter with schemas.enable=true reading bare JSON; or a key written by StringConverter but read by AvroConverter. For resilience, configure errors.tolerance=all with a dead-letter queue (errors.deadletterqueue.topic.name, plus errors.deadletterqueue.context.headers.enable=true to capture the cause) so poison records are routed aside instead of stopping the task. Also verify schema-registry connectivity/auth and compatibility settings. Fix by correcting the converter config or republishing data in the expected format.
code
properties · 6 lineserrors.tolerance=all
errors.deadletterqueue.topic.name=dlq-myconnector
errors.deadletterqueue.topic.replication.factor=3
errors.deadletterqueue.context.headers.enable=true
errors.log.enable=true
errors.log.include.messages=truego deeper
Recognize that the error usually means the sink's converter doesn't match how the data was written.
Inspect the on-topic format, align key/value converter settings, and know schemas.enable/schema.registry.url must match.
Drive a systematic diagnosis and configure errors.tolerance + DLQ with context headers for resilience.
Establish org-wide converter/registry standards and contracts so producer/consumer format drift can't happen.
## The root cause A sink converter deserializes topic bytes into Connect data. If the configured converter does not match how the data was actually serialized, you get a `DataException` (Connect) wrapping a `SerializationException` or a registry error. **The data and the converter disagree.** ## Step 1 — Determine the actual on-topic format Consume raw bytes (e.g. `kafka-console-consumer` without a value deserializer, or inspect with a hex view) and look at the first bytes: - Leading `0x00` magic byte + 4-byte ID -> Confluent registry format (Avro/Protobuf/JsonSchema). The sink must use the matching registry-backed converter and `schema.registry.url`. - Leading `{` with `"schema"`/`"payload"` -> JsonConverter with `schemas.enable=true`. - Leading `{` bare object -> JsonConverter with `schemas.enable=false`. - Plain text -> StringConverter. - Arbitrary binary, connector-owned -> ByteArrayConverter. Check the **key** separately from the value — a frequent trap is the value being right but the key written with a different converter (e.g. String key, Avro value). ## Step 2 — Align converter config Set the sink's converters and their sub-properties to match: ``` value.converter=io.confluent.connect.avro.AvroConverter value.converter.schema.registry.url=http://sr:8081 key.converter=org.apache.kafka.connect.storage.StringConverter ``` For JSON, match `schemas.enable` on both producer and sink. For registry converters, ensure the schema ID in the message resolves in the registry the sink is pointed at (wrong/empty registry, auth failure, or network partition also throws here). ## Step 3 — Common concrete failures - **AvroConverter sink on JSON data:** the converter reads the magic byte, gets a nonsensical schema ID, and throws. Fix: use JsonConverter, or republish as Avro. - **schemas.enable mismatch:** `true` sink reading bare JSON tries to find `schema`/`payload` and fails; `false` sink reading the envelope produces a wrong/nested structure. - **Schema compatibility:** registry rejects or the reader can't evolve to the writer schema (rare on read for Avro with proper compatibility, but possible with strict reader schemas / `use.latest.version`). ## Step 4 — Make the pipeline resilient (error handling) Rather than letting one poison record stop the task: ``` errors.tolerance=all errors.deadletterqueue.topic.name=dlq-myconnector errors.deadletterqueue.context.headers.enable=true errors.log.enable=true errors.log.include.messages=true ``` - `errors.tolerance=all` skips records that fail in conversion/SMT/put instead of failing the task (`none` is the default and fails fast). - The **dead-letter queue** routes failed records to a separate topic; with `context.headers.enable=true`, headers on the DLQ record carry the exception class, message, and original topic/partition/offset for forensics. - `errors.log.*` writes the failures to the Connect log. Note: the DLQ applies to **sink** connectors. Converter and SMT errors are covered by this framework; errors thrown deep inside the connector's external write may not all be. ## Step 5 — Decide the durable fix Either correct the converter configuration (the usual fix) or, if upstream is wrong, republish/standardize the data format. Long-term, standardize converters and schema-registry usage across the org so producers and consumers can't drift.
- How do you keep one bad record from stopping the whole sink task?Set errors.tolerance=all and configure a dead-letter queue (errors.deadletterqueue.topic.name) with errors.deadletterqueue.context.headers.enable=true so poison records are routed aside with their failure cause instead of failing the task.
- How can you tell from the bytes whether a topic uses a registry-backed converter?Registry-backed (Avro/Protobuf/JsonSchema) messages begin with a 0x00 magic byte followed by a 4-byte big-endian schema ID. Plain JSON or text won't have that prefix.
- Why must you check the key converter even if the value deserializes fine?Key and value are converted independently; a common mismatch is a String/Int key being read by a registry-backed key converter (or vice versa), which fails even when the value is correct.
saying these in an interview costs you the question
- Blaming the connector code when it's a converter/format mismatch
- Forgetting that errors.tolerance defaults to none (fail-fast)
- Assuming the DLQ captures the exception cause without enabling context.headers
- Ignoring the key converter and only checking the value
- Thinking a registry-backed sink can read plain JSON because 'it's all just JSON'