【问题标题】: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();
}