{"record":{"id":"680d54e225f466e9","repo":"apache/beam","slug":"upgrading-kafkaio-write-transforms-that-have","errorCode":null,"errorMessage":"Upgrading KafkaIO write transforms that have `withBadRecordErrorHandler` property set is not supported yet.","messagePattern":"Upgrading KafkaIO write 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":580,"sourceCode":"      if (writeRecordsTransform.getConsumerFactoryFn() != null) {\n        fieldValues.put(\n            \"consumer_factory_fn\", toByteArray(writeRecordsTransform.getConsumerFactoryFn()));\n      }\n\n      if (writeRecordsTransform.getProducerConfig().size() > 0) {\n        Map<String, byte[]> producerConfigMap = new HashMap<>();\n        writeRecordsTransform\n            .getProducerConfig()\n            .forEach(\n                (key, value) -> {\n                  producerConfigMap.put((String) key, toByteArray(value));\n                });\n        fieldValues.put(\"producer_config\", producerConfigMap);\n      }\n      if (writeRecordsTransform.getBadRecordErrorHandler() != null\n          && !(writeRecordsTransform.getBadRecordErrorHandler()\n              instanceof ErrorHandler.DefaultErrorHandler)) {\n        throw new RuntimeException(\n            \"Upgrading KafkaIO write transforms that have `withBadRecordErrorHandler` property set\"\n                + \" is not supported yet.\");\n      }\n\n      return Row.withSchema(schema).withFieldValues(fieldValues).build();\n    }\n\n    @Override\n    public Write<?, ?> fromConfigRow(Row configRow, PipelineOptions options) {\n      try {\n        Write<?, ?> transform = KafkaIO.write();\n\n        String bootstrapServers = configRow.getString(\"bootstrap_servers\");\n        if (bootstrapServers != null) {\n          transform = transform.withBootstrapServers(bootstrapServers);\n        }\n        String topic = configRow.getValue(\"topic\");\n        if (topic != null) {","sourceCodeStart":562,"sourceCodeEnd":598,"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#L562-L598","documentation":"During Beam pipeline upgrade/translation, KafkaIO write transforms are converted to a config Row via toConfigRow. If the transform has a custom bad-record error handler attached (withBadRecordErrorHandler) that is not the default handler, the upgrade path cannot represent it, so translation fails with this RuntimeException.","triggerScenarios":"Running pipeline upgrade (toConfigRow/row) on a KafkaIO<KafkaRecord<Void,V>>.writeRecords() transform where .withBadRecordErrorHandler(handler) was called with a custom ErrorHandler implementation (anything other than ErrorHandler.DefaultErrorHandler).","commonSituations":"Migrating an existing Beam Kafka sink pipeline to the newer KafkaIO managed/upgrade representation when the pipeline was written with custom dead-letter error handling for producer failures.","solutions":["Remove the custom withBadRecordErrorHandler from the KafkaIO write transform before upgrading, or temporarily use the default ErrorHandler.DefaultErrorHandler","Implement equivalent bad-record handling downstream of the sink (e.g., a separate dead-letter write path) instead of via withBadRecordErrorHandler","Wait for / upgrade to a Beam version that supports translating withBadRecordErrorHandler for KafkaIO writes"],"exampleFix":"// before\nKafkaIO.<KafkaRecord<Void, byte[]>>writeRecords()\n    .withBootstrapServers(servers)\n    .withTopic(topic)\n    .withBadRecordErrorHandler(customErrorHandler);\n// after\nKafkaIO.<KafkaRecord<Void, byte[]>>writeRecords()\n    .withBootstrapServers(servers)\n    .withTopic(topic);","handlingStrategy":"validation","validationCode":"if (transform.getBadRecordErrorHandler() != null\n    && !(transform.getBadRecordErrorHandler() instanceof ErrorHandler.DefaultErrorHandler)) {\n  throw new IllegalStateException(\"Detach custom bad-record handler before upgrading KafkaIO write transform\");\n}","typeGuard":null,"tryCatchPattern":"try {\n  row = translation.toConfigRow(transform);\n} catch (RuntimeException e) {\n  if (e.getMessage().contains(\"withBadRecordErrorHandler\")) {\n    // re-plan without custom handler or abort upgrade with clear user guidance\n  }\n}","preventionTips":["Avoid withBadRecordErrorHandler with custom handlers on KafkaIO writes if you plan to use the upgrade/managed path","Keep dead-letter handling as a separate downstream transform","Pin Beam version and check release notes for translation support before upgrading"],"tags":["java","apache-beam","kafka","pipeline-upgrade","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"}