{"record":{"id":"7692a8f888f22c61","repo":"apache/beam","slug":"writetokafka-with-headers-true-only-supports","errorCode":null,"errorMessage":"WriteToKafka(with_headers=True) only supports ByteArraySerializer for key and value.","messagePattern":"WriteToKafka\\(with_headers=True\\) only supports ByteArraySerializer for key and value\\.","errorType":"validation","errorClass":"ValueError","httpStatus":null,"severity":"error","filePath":"sdks/python/apache_beam/io/kafka.py","lineNumber":325,"sourceCode":"\n    :param producer_config: A dictionary containing the producer configuration.\n    :param topic: A Kafka topic name.\n    :param key_serializer: A fully-qualified Java class name of a Kafka\n        Serializer for the topic's key, e.g.\n        'org.apache.kafka.common.serialization.LongSerializer'.\n        Default: 'org.apache.kafka.common.serialization.ByteArraySerializer'.\n    :param value_serializer: A fully-qualified Java class name of a Kafka\n        Serializer for the topic's value, e.g.\n        'org.apache.kafka.common.serialization.LongSerializer'.\n        Default: 'org.apache.kafka.common.serialization.ByteArraySerializer'.\n    :param with_headers: If True, input elements must be beam.Row objects\n        containing 'key', 'value', and optional 'headers' fields.\n        Only ByteArraySerializer is supported when with_headers=True.\n    :param expansion_service: The address (host:port) of the ExpansionService.\n    \"\"\"\n    if with_headers and (key_serializer != self.byte_array_serializer or\n                         value_serializer != self.byte_array_serializer):\n      raise ValueError(\n          'WriteToKafka(with_headers=True) only supports '\n          'ByteArraySerializer for key and value.')\n\n    urn = self.URN_WITH_HEADERS if with_headers else self.URN\n    super().__init__(\n        urn,\n        NamedTupleBasedPayloadBuilder(\n            WriteToKafkaSchema(\n                producer_config=producer_config,\n                topic=topic,\n                key_serializer=key_serializer,\n                value_serializer=value_serializer,\n            )),\n        expansion_service or default_io_expansion_service())\n","sourceCodeStart":307,"sourceCodeEnd":340,"githubUrl":"https://github.com/apache/beam/blob/12126d8942aaf848030c478b4c6a28c6af861c66/sdks/python/apache_beam/io/kafka.py#L307-L340","documentation":"WriteToKafka with with_headers=True uses the external Kafka expansion service's header-bearing transform, which requires keys and values to be raw bytes serialized with Kafka's ByteArraySerializer. If key_serializer or value_serializer is anything else, the constructor raises ValueError because the external transform cannot accept other serializers in headers mode.","triggerScenarios":"WriteToKafka(..., with_headers=True, key_serializer='org.apache.kafka.common.serialization.StringSerializer', ...) or any serializer besides WriteToKafka.byte_array_serializer ('org.apache.kafka.common.serialization.ByteArraySerializer').","commonSituations":"Needing per-record headers (common for tracing/audit metadata) while keeping default String serializers; copying an existing WriteToKafka call and only adding with_headers=True.","solutions":["Set key_serializer and value_serializer to WriteToKafka.byte_array_serializer and encode keys/values to bytes yourself.","Disable with_headers if you don't need headers and keep your existing serializers.","Pre-encode records (e.g. JSON/Avro to bytes) before writing when switching to ByteArraySerializer."],"exampleFix":"# before\nWriteToKafka(bootstrap_servers, topic, with_headers=True)\n# after\nWriteToKafka(bootstrap_servers, topic, with_headers=True,\n             key_serializer=WriteToKafka.byte_array_serializer,\n             value_serializer=WriteToKafka.byte_array_serializer)","handlingStrategy":"validation","validationCode":"if with_headers and not (key_serializer == WriteToKafka.byte_array_serializer and\n                          value_serializer == WriteToKafka.byte_array_serializer):\n    raise ValueError('with_headers=True requires ByteArraySerializer')","typeGuard":"def supports_headers(key_serializer, value_serializer) -> bool:\n    ba = WriteToKafka.byte_array_serializer\n    return key_serializer == ba and value_serializer == ba","tryCatchPattern":"try:\n    _ = WriteToKafka(bootstrap, topic, with_headers=True, key_serializer=ks, value_serializer=vs)\nexcept ValueError as e:\n    if 'ByteArraySerializer' in str(e): log.error('Switch serializers or disable with_headers')","preventionTips":["Encode keys/values to bytes before writing when using headers","Only add with_headers=True after switching serializers","Check serializer names against WriteToKafka.byte_array_serializer"],"tags":["apache-beam","python","kafka","unsupported-serializer","headers"],"backgroundTag":"conflicting-config-options","analyzedSha":"12126d8942aaf848030c478b4c6a28c6af861c66","analyzedAt":"2026-09-13T01:50:10.254Z","contentChangedAt":"2026-09-13T01:50:10.254Z","schemaVersion":2},"datasetVersion":"2026-09-20T03:17:13.778Z"}