njL*_*jLT 5 java google-cloud-pubsub apache-beam apache-beam-io
我正在尝试使用以下代码从 pub/sub 中读取
Read<String> pubsub = PubsubIO.<String>read().topic("projects/<projectId>/topics/<topic>").subscription("projects/<projectId>/subscriptions/<subscription>").withCoder(StringUtf8Coder.of()).withAttributes(new SimpleFunction<PubsubMessage,String>() {
@Override
public String apply(PubsubMessage input) {
LOG.info("hola " + input.getAttributeMap());
return new String(input.getMessage());
}
});
PCollection<String> pps = p.apply(pubsub)
.apply(
Window.<String>into(
FixedWindows.of(Duration.standardSeconds(15))));
pps.apply("printdata",ParDo.of(new DoFn<String, String>() {
@ProcessElement
public void processElement(ProcessContext c) {
LOG.info("hola amigo "+c.element());
c.output(c.element());
}
}));
Run Code Online (Sandbox Code Playgroud)
与我在 NodeJS 上收到的相比,我收到了将包含在该data字段中的消息。我怎样才能得到这个ackId字段(我以后可以用它来确认消息)?我正在打印的属性映射是null. 是否有其他方法可以确认所有消息而无需弄清楚 ackId?
该PubsubIO读者负责确认消息。它与跑步者的检查点行为有关。具体来说,源只会在结果元素被检查点时确认消息。
在这种情况下,您应该查看 Flink 运行程序检查点何时提供该源的状态信息。我相信这与检查点频率的 Flink 配置有关。
| 归档时间: |
|
| 查看次数: |
1494 次 |
| 最近记录: |