Acknowledge Google Pub/Sub message on Apache Beam

Viewed 1822

I'm trying to read from pub/sub with the following code

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());
    }
  }));

Compared to what I receive on NodeJS, I get the message that would be contained in the data field. How can I get the ackId field (which I can later use to acknowledge the message)? The attribute map that I'm printing is null. Is there some other way to acknowledge all messages without having to figure out the ackId?

1 Answers
Related