【问题标题】:Dataflow: string to pubsub message数据流:发布订阅消息的字符串
【发布时间】:2019-12-17 13:54:33
【问题描述】:

我正在尝试在 Dataflow 中进行单元测试。

对于那个测试,在乞讨时,我将从一个简单的硬编码字符串开始。

问题是我需要将该字符串转换为 pubsub 消息。我得到了以下代码:

    // Create a PCollection from string a transform to pubsub message format
    PCollection<PubsubMessage> input = p.apply("input string", Create.of("test" + 
            ""))
            .apply("convert to Pub/Sub message", ParDo.of(new DoFn<String, PubsubMessage>() {
                @ProcessElement
                public void processElement(ProcessContext c) {
                    c.output(new PubsubMessage(c.element().getBytes(), null));
                }
            }));

但我收到以下错误:

 java.lang.IllegalArgumentException: unable to serialize DoFnWithExecutionInformation{doFn=com.xxx.pipeline.TesterPipeline$1@7b64240d, mainOutputTag=Tag<output>, sideInputMapping={}, schemaInformation=DoFnSchemaInformation{elementConverters=[]}}
    at org.apache.beam.sdk.util.SerializableUtils.serializeToByteArray(SerializableUtils.java:55)
    <...>
Caused by: java.io.NotSerializableException: com.xxx.pipeline.TesterPipeline
    at java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1184)
    at java.io.ObjectOutputStream.defaultWriteFields(ObjectOutputStream.java:1548)
    at java.io.ObjectOutputStream.writeSerialData(ObjectOutputStream.java:1509)
    at java.io.ObjectOutputStream.writeOrdinaryObject(ObjectOutputStream.java:1432)
    at java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1178)
    at java.io.ObjectOutputStream.defaultWriteFields(ObjectOutputStream.java:1548)
    at java.io.ObjectOutputStream.writeSerialData(ObjectOutputStream.java:1509)
    at java.io.ObjectOutputStream.writeOrdinaryObject(ObjectOutputStream.java:1432)
    at java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1178)
    at java.io.ObjectOutputStream.writeObject(ObjectOutputStream.java:348)
    at org.apache.beam.sdk.util.SerializableUtils.serializeToByteArray(SerializableUtils.java:51)
    ... 50 more

我应该如何从字符串创建 pubsub 消息?

【问题讨论】:

    标签: java google-cloud-dataflow apache-beam


    【解决方案1】:

    在用户 ParDos 的 Serializability 要求下的 Beam Programming Guide 中,它提到了这一点:

    使用匿名内部类实例来内联声明函数对象时要小心。在非静态上下文中,您的内部类实例将隐式包含指向封闭类和该类状态的指针。该封闭类也将被序列化,因此适用于函数对象本身的相同注意事项也适用于这个外部类。

    发生的情况是您的匿名 DoFn 隐式包含指向您正在构建管道的类的指针,这导致了此序列化失败。您可以通过将 DoFn 设为命名子类而不是匿名来避免这种情况:

    public class MyDoFn extends DoFn<String, PubsubMessage>() {
      @ProcessElement
      public void processElement(ProcessContext c) {
        c.output(new PubsubMessage(c.element().getBytes(), null));
      }
    }
    

    【讨论】:

      【解决方案2】:
      猜你喜欢
      • 2021-11-30
      • 1970-01-01
      • 1970-01-01
      • 2016-07-03
      • 2016-11-25
      • 2021-02-27
      • 2015-06-30
      • 2015-03-01
      • 2016-04-18
      相关资源
      最近更新 更多