{"record":{"id":"2adf25edade6197d","repo":"apache/beam","slug":"the-current-runner-does-not-support-sdf-based-kafka-read","errorCode":null,"errorMessage":"The current runner does not support SDF-based Kafka read properly and the replacement runner lacks the support for the following properties: %s. For example if you are using Dataflow then consider using Dataflow Runner v2.","messagePattern":"The current runner does not support SDF-based Kafka read properly and the replacement runner lacks the support for the following properties: (.+?)\\. For example if you are using Dataflow then consider using Dataflow Runner v2\\.","errorType":"exception","errorClass":"IllegalStateException","httpStatus":null,"severity":"error","filePath":"sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java","lineNumber":1875,"sourceCode":"            new KafkaReadOverrideFactory<>());\n\n    private static class KafkaReadOverrideFactory<K, V>\n        implements PTransformOverrideFactory<\n            PBegin, PCollection<KafkaRecord<K, V>>, ReadFromKafkaViaSDF<K, V>> {\n\n      @Override\n      public PTransformReplacement<PBegin, PCollection<KafkaRecord<K, V>>> getReplacementTransform(\n          AppliedPTransform<PBegin, PCollection<KafkaRecord<K, V>>, ReadFromKafkaViaSDF<K, V>>\n              transform) {\n        try {\n          return PTransformReplacement.of(\n              transform.getPipeline().begin(),\n              new ReadFromKafkaViaUnbounded<>(\n                  transform.getTransform().kafkaRead,\n                  transform.getTransform().keyCoder,\n                  transform.getTransform().valueCoder));\n        } catch (KafkaIOReadImplementationCompatibilityException e) {\n          throw new IllegalStateException(\n              \"The current runner does not support SDF-based Kafka read properly \"\n                  + \"and the replacement runner lacks the support for the following properties: \"\n                  + e.getConflictingProperties()\n                  + \". For example if you are using Dataflow then consider using Dataflow Runner v2.\");\n        }\n      }\n\n      @Override\n      public Map<PCollection<?>, ReplacementOutput> mapOutputs(\n          Map<TupleTag<?>, PCollection<?>> outputs, PCollection<KafkaRecord<K, V>> newOutput) {\n        return ReplacementOutputs.singleton(outputs, newOutput);\n      }\n    }\n\n    private abstract static class AbstractReadFromKafka<K, V>\n        extends PTransform<PBegin, PCollection<KafkaRecord<K, V>>> {\n      Read<K, V> kafkaRead;\n      Coder<K> keyCoder;","sourceCodeStart":1857,"sourceCodeEnd":1893,"githubUrl":"https://github.com/apache/beam/blob/12126d8942aaf848030c478b4c6a28c6af861c66/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java#L1857-L1893","documentation":"When expanding KafkaIO.Read, Beam tries to replace the legacy read with an SDF-based implementation (ReadFromKafkaViaUnbounded). If the runner does not support SDF or cannot honor required properties (KafkaIOReadImplementationCompatibilityException), expansion fails. Dataflow v1 is the classic case; Runner v2 supports SDF.","triggerScenarios":"Expanding KafkaIO.read() on a runner whose SDF support lacks properties reported by the compatibility check — e.g. Dataflow legacy runner, or a runner without watermark/track-property support — during pipeline construction.","commonSituations":"Running on old Dataflow (Runner v1), an outdated runner version without full SDF support, or a custom/portable runner missing required SDF capabilities.","solutions":["Switch to Dataflow Runner v2 (--experiments=beam_fnapi / use Runner v2 default) if using Dataflow.","Upgrade the runner to a version with full SDF support.","Run on DirectRunner/Flink/Spark which support the SDF-based read, for local testing.","Set --useDeprecatedKafkaRead (if available in your Beam version) to keep the legacy read implementation."],"exampleFix":"// before (Dataflow v1)\n--runner=DataflowRunner\n// after\n--runner=DataflowRunner --experiments=use_runner_v2","handlingStrategy":"validation","validationCode":"// Before submitting: verify runner supports SDF\nif (options.getRunner() == DataflowRunner.class && !options.getExperiments().contains(\"use_runner_v2\")) {\n  options.setExperiments(ImmutableList.of(\"use_runner_v2\"));\n}","typeGuard":null,"tryCatchPattern":"try {\n  pipeline.run();\n} catch (IllegalStateException e) {\n  if (e.getMessage() != null && e.getMessage().contains(\"SDF-based Kafka read\")) {\n    throw new IllegalStateException(\"Re-run with Dataflow Runner v2: --experiments=use_runner_v2\", e);\n  }\n  throw e;\n}","preventionTips":["Use Dataflow Runner v2 (beam_fnapi) for Kafka IO","Keep runner versions current with your Beam SDK version","Smoke-test Kafka reads with DirectRunner before switching runners"],"tags":["java","kafka","runner","sdf","dataflow"],"backgroundTag":"unsupported-operation","analyzedSha":"12126d8942aaf848030c478b4c6a28c6af861c66","analyzedAt":"2026-09-13T01:50:10.254Z","contentChangedAt":"2026-09-13T01:50:10.254Z","schemaVersion":2},"datasetVersion":"2026-09-14T16:17:12.679Z"}