{"record":{"id":"6da61b10cfb815e8","repo":"apache/beam","slug":"upgrading-kafkaio-read-transforms-that-have","errorCode":null,"errorMessage":"Upgrading KafkaIO read transforms that have `withBadRecordErrorHandler` property set is not supported yet.","messagePattern":"Upgrading KafkaIO read transforms that have `withBadRecordErrorHandler` property set is not supported yet\\.","errorType":"exception","errorClass":"RuntimeException","httpStatus":null,"severity":"error","filePath":"sdks/java/io/kafka/upgrade/src/main/java/org/apache/beam/sdk/io/kafka/upgrade/KafkaIOTranslation.java","lineNumber":219,"sourceCode":"            .forEach(\n                (key, val) -> {\n                  offsetConsumerConfigMap.put(key, toByteArray(val));\n                });\n        fieldValues.put(\"offset_consumer_config\", offsetConsumerConfigMap);\n      }\n      if (transform.getKeyDeserializerProvider() != null) {\n        fieldValues.put(\n            \"key_deserializer_provider\", toByteArray(transform.getKeyDeserializerProvider()));\n      }\n      if (transform.getValueDeserializerProvider() != null) {\n        fieldValues.put(\n            \"value_deserializer_provider\", toByteArray(transform.getValueDeserializerProvider()));\n      }\n      if (transform.getCheckStopReadingFn() != null) {\n        fieldValues.put(\"check_stop_reading_fn\", toByteArray(transform.getCheckStopReadingFn()));\n      }\n      if (transform.getBadRecordErrorHandler() != null) {\n        throw new RuntimeException(\n            \"Upgrading KafkaIO read transforms that have `withBadRecordErrorHandler` property set\"\n                + \" is not supported yet.\");\n      }\n      if (transform.getLogTopicVerification() != null) {\n        fieldValues.put(\"log_topic_verification\", transform.getLogTopicVerification());\n      }\n\n      fieldValues.put(\"redistribute\", transform.isRedistributed());\n      fieldValues.put(\"redistribute_num_keys\", transform.getRedistributeNumKeys());\n      fieldValues.put(\"allows_duplicates\", transform.isAllowDuplicates());\n      if (transform.getOffsetDeduplication() != null) {\n        fieldValues.put(\"offset_deduplication\", transform.getOffsetDeduplication());\n      }\n      if (transform.getRedistributeByRecordKey() != null) {\n        fieldValues.put(\"redistribute_by_record_key\", transform.getRedistributeByRecordKey());\n      }\n      return Row.withSchema(schema).withFieldValues(fieldValues).build();\n    }","sourceCodeStart":201,"sourceCodeEnd":237,"githubUrl":"https://github.com/apache/beam/blob/12126d8942aaf848030c478b4c6a28c6af861c66/sdks/java/io/kafka/upgrade/src/main/java/org/apache/beam/sdk/io/kafka/upgrade/KafkaIOTranslation.java#L201-L237","documentation":"During pipeline upgrade, KafkaIOTranslation.toConfigRow serializes a KafkaIO read transform into a config Row. The `withBadRecordErrorHandler` property has no serialized representation yet, so if the transform sets it, the translator refuses to encode the transform and throws, preventing the pipeline from being upgraded/portable-ized.","triggerScenarios":"Calling KafkaIO.Read.create()...withBadRecordErrorHandler(handler) and then running the pipeline through the upgrade path (row()/toConfigRow payload translation, e.g. legacy-to-portable upgrade or ReadRegistrar payload translation).","commonSituations":"Pipelines that use bad-record error handling on Kafka reads being migrated between Beam versions or runner environments; a developer added withBadRecordErrorHandler to an existing pipeline and then tried to upgrade/translate it.","solutions":["Remove the withBadRecordErrorHandler call from the transform before upgrading.","Handle bad records manually (e.g. wrap the downstream DoFn with error capture, or use a separate DLQ sink) instead of relying on the built-in handler.","Upgrade the pipeline on a Beam version where the handler is unset, then re-add the property after translation support lands.","Track Beam's KafkaIOTranslation changes; once supported, re-add withBadRecordErrorHandler."],"exampleFix":"// before\nKafkaIO.<byte[], byte[]>read().withBadRecordErrorHandler(handler)\n// after\nKafkaIO.<byte[], byte[]>read() // no bad record handler during upgrade\n","handlingStrategy":"validation","validationCode":"if (transform.getBadRecordErrorHandler() != null) {\n  throw new IllegalStateException(\"Remove withBadRecordErrorHandler before pipeline upgrade/translation.\");\n}","typeGuard":null,"tryCatchPattern":null,"preventionTips":["Avoid withBadRecordErrorHandler on Kafka read transforms in pipelines subject to legacy-to-portable upgrades.","Centralize bad-record handling in downstream DoFns/DLQ sinks instead of the transform property.","Check the Beam version's upgrade support notes before adding new KafkaIO properties."],"tags":["kafka","pipeline-upgrade","serialization","unsupported-feature"],"backgroundTag":"unsupported-operation","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"}