Skip to content
Back to Insights
Data EngineeringBy KE Engineering Team

KACC

KACCDATA ENGINEERING cover for KACCTOPICc0c1c2groupDATA ENGINEERINGKACC// BYTES OVER STRINGS

First written September 2022, last updated September 2026.

Introduction

The Confluent Avro serializer and deserializer store the unique ID of the schema in the message. When unexpected characters show up in a string, a type mismatch is obvious. But what about non-printable characters? How do they show up? Will the issue then be obvious?

Demonstration

A simple demonstration can be done with the Datagen Source Connector. Create a connector with Avro as the key. The data type for the Datagen’s quickstart users is a string. The Avro serializer writes this as an Avro primitive. Typically, when Avro is used, the top-level object is a Record, but the serializer has custom code for supporting primitives.

The Configuration

The Datagen connector is configured with the key represented as Avro.

json
{    "connector.class": "io.confluent.kafka.connect.datagen.DatagenConnector",    "tasks.max": "1",    "kafka.topic": "users",    "quickstart": "users",    "key.converter": "io.confluent.connect.avro.AvroConverter",    "key.converter.schema.registry.url" : "http://schema-registry:8081",    "key.converter.schemas.enable": "true",    "value.converter": "io.confluent.connect.avro.AvroConverter",    "value.converter.schema.registry.url" : "http://schema-registry:8081",    "value.converter.schemas.enable": "true",    "max.interval": 100,    "iterations": 10000000}

Scenario

You write a Kafka Streams application where you read the key as a Serdes.String(), the default you used for your application. You forget to change the serde for reading users from the default serde to an Avro Serde. You now join your stream of orders with users, and none of the joins succeeds.

Investigation…

The first thing most of us do is reach for kafka-avro-console-consumer to see what is going on.

shell
kafka-avro-console-consumer \        --bootstrap-server localhost:19092 \        --property schema.registry.url="http://localhost:8081" \        --property print.key=true \        --property key.separator="|" \        --from-beginning \        --skip-message-on-error \        --key-deserializer=org.apache.kafka.common.serialization.StringDeserializer \        --topic users

The result has content that looks pretty normal and expected:

text
User_9|{"registertime":1489457902486,"userid":"User_9","regionid":"Region_1","gender":"OTHER"}User_1|{"registertime":1500277798184,"userid":"User_1","regionid":"Region_2","gender":"OTHER"}

There could be extra blank lines if the non-printable bytes trigger them, but that doesn’t always stand out as an obvious issue.

If your key deserializer were BytesDeserializer, what would you have seen?

shell
kafka-avro-console-consumer \        --bootstrap-server localhost:19092 \        --property schema.registry.url="http://localhost:8081" \        --property print.key=true \        --property key.separator="|" \        --from-beginning \        --skip-message-on-error \        --key-deserializer=org.apache.kafka.common.serialization.BytesDeserializer \        --topic users

The serializer’s magic byte (0x00), the four bytes of the schema id (3), and Avro's zigzag-encoded length prefix (0x0C decodes to 6, the length of User_9) show up in printable hex characters:

text
\x00\x00\x00\x00\x03\x0CUser_9|{"registertime":1489457902486,"userid":"User_9","regionid":"Region_1","gender":"OTHER"}\x00\x00\x00\x00\x03\x0CUser_1|{"registertime":1500277798184,"userid":"User_1","regionid":"Region_2","gender":"OTHER"}

Now it is easy to see the issue: the key is Avro (a primitive Avro string as defined by the serializer). Solution: change the connector to write String keys, or read the key with the Avro serde in your Streams app.

NOTE

Running containers for demonstrations is great, but the mismatch of URLs can be confusing. localhost:port is used for connecting to services from the host machine (your laptop) via port mapping. The actual hostname is used when you are accessing the service from another container. Therefore, you will see http://schema-registry:8081 within the connect configuration, and http://localhost:8081 for commands running from the host machine. We haven't translated them here, since these scripts align with the demo code.

Useful Shell Aliases

We have these defined in our .zshrc.

shell
alias kcc='kafka-console-consumer \        --bootstrap-server localhost:19092 \        --key-deserializer=org.apache.kafka.common.serialization.BytesDeserializer  \        --property print.key=true \        --property key.separator="|" \        --from-beginning \        --topic'
shell
alias kacc='kafka-avro-console-consumer \        --bootstrap-server localhost:19092 \        --property schema.registry.url="http://localhost:8081" \        --property print.key=true \        --property key.separator="|" \        --from-beginning \        --skip-message-on-error \        --key-deserializer=org.apache.kafka.common.serialization.BytesDeserializer \        --topic'

Takeaways

  • This may seem obvious, since you'd inspect the connector configuration and find the problem immediately, but you want it to be easy for everyone on your team.
  • This demonstration is available as the key-mismatch demo in dev-local-demos.

Working on something like this?

Start a Conversation