Apache Beam中如何为PCollection<List<String>>类型设置Coder
问题根因
你遇到的错误核心有两个:一是Java泛型擦除导致无法直接获取List<String>的类型信息推断Coder,二是自定义Transform里中间输出的JsonNode类型PCollection也没有指定Coder,两个问题同时存在才会出现设置了List的Coder仍然报错的情况。
正确解决方案
方案1:直接在PCollection后指定Coder(最稳妥)
你可以任选位置添加,覆盖所有缺Coder的节点:
- 改自定义Transform内部的
expand方法,同时处理中间节点和输出节点:
@Override public PCollection<List<String>> expand(PCollection<String> input) { return input .apply(MapElements.via(new GetRootNode())) // 给JsonNode中间结果加Coder .setCoder(SerializableCoder.of(JsonNode.class)) .apply(MapElements.via(new ExtractPathsFromTree())) // 给最终的List<String>结果加Coder,内置StringUtf8Coder比SerializableCoder性能更优 .setCoder(ListCoder.of(StringUtf8Coder.of())); }
- 不想改Transform代码的话,直接在流水线调用时给输出加(内部的JsonNode节点Coder还是要在Transform里补):
.apply("Traverse Json tree", new JSONTreeToPaths()) .setCoder(ListCoder.of(StringUtf8Coder.of()))
方案2:通过TypeDescriptor声明输出类型
定义MapElements时直接声明输出类型,Beam会自动匹配默认Coder,不需要手动调用setCoder:
.apply(MapElements.via(new ExtractPathsFromTree()) .withOutputType(TypeDescriptor.of(new TypeLiteral<List<String>>() {})))
JsonNode的转换步骤也可以用同样方式声明类型。
之前尝试失败的原因
SerializableCoder.of(List<String>.class)报错是因为Java泛型擦除,运行时不存在List<String>.class这个对象,只有List.class,无法读取泛型参数类型。- 用
ListCoder.of(SerializableCoder.of(String.class))仍然报错,是因为你只处理了List类型的输出Coder,没补全中间JsonNode节点的Coder,报错实际上是中间节点抛出的。
内容的提问来源于stack exchange,提问作者rocksNwaves
相关产品推荐
相关产品推荐

