【问题标题】:Reading from Pubsub using Dataflow Java SDK 2使用 Dataflow Java SDK 2 从 Pubsub 读取
【发布时间】:2018-01-22 16:44:16
【问题描述】:

Google Cloud Platform for Java SDK 2.x 的许多文档都告诉您参考 Beam 文档。

当使用 Dataflow 从 PubSub 读取数据时,我是否仍然在执行 PubsubIO.Read.named("name").topic("");

或者我应该做点别的吗?

此外,有没有办法将 Dataflow 接收到的 PubSub 数据打印到标准输出或文件?

【问题讨论】:

    标签: google-cloud-platform google-cloud-dataflow apache-beam google-cloud-pubsub


    【解决方案1】:

    对于 Apache Beam 2.2.0,您可以定义以下转换以从 Pub/Sub 订阅中提取消息:

    PubsubIO.readMessages().fromSubscription("subscription_name")
    

    这是定义将从 Pub/Sub 中提取消息的转换的一种方法。但是,PubsubIO 类包含用于拉取消息的不同方法。每种方法的功能略有不同。请参阅PubsubIO 文档。

    您可以使用 TextIO 类将 Pub/Sub 消息写入文件。请参阅TextIO 文档中的示例。有关将 Pub/Sub 消息写入 stdout 的信息,请参阅 Logging Pipeline Messages 文档。

    【讨论】:

      【解决方案2】:

      添加到 Adrew 上面写的内容。从 PubSubIO 读取字符串并将它们写入标准输出(仅用于调试)的代码如下。也就是说,我将提交内部错误以改进 PubsubIO 的 JavaDoc,我认为当前的文档很少。

      public static void main(String[] args) {
      
        Pipeline pipeline = Pipeline.create(PipelineOptionsFactory.fromArgs(args).create());
        pipeline
          .apply("ReadStrinsFromPubsub",
             PubsubIO.readStrings().fromTopic("/topics/my_project/my_topic"))
          .apply("PrintToStdout", ParDo.of(new DoFn<String, Void>() {
            @ProcessElement
            public void processElement(ProcessContext c) {
              System.out.printf("Received at %s : %s\n", Instant.now(), c.element()); // debug log
            }
          }));
      
        pipeline.run().waitUntilFinish();
      }
      

      【讨论】:

      • 我们需要先注册订阅者还是可以直接阅读主题?
      猜你喜欢
      • 2018-05-13
      • 2020-07-27
      • 1970-01-01
      • 2020-10-16
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多