Apache Beam pipeline with FlinkRunner does not send ack's on PubSub when reading messages

Viewed 213

I have a pipeline that reads messages using PubsubIO, this pipeline publishes a lot of messages, but unfortunately I don't get the acknowledgment. This is the code that reads the messages

public CampaignMetricStream() {
        super(Options.class);
    }

    @Override
    public PCollection<String> source(Pipeline pipeline, Options options) {
        return pipeline.apply("pubsub-pull-events", PubsubIO
                .readStrings()
                .fromSubscription(options.getSubscription()));
    }

I'm pretty sure that it is not a problem in the code, but in the configuration.

Flink is configured on a k8s cluster

    taskmanager.numberOfTaskSlots: 64
    taskmanager.memory.managed.size: 0
    jobmanager.memory.process.size: 1g
    taskmanager.memory.process.size: 50g
    parallelism.default: 2
    state.backend: filesystem
    state.checkpoints.dir: file:///data/flink/checkpoints
    taskmanager.memory.jvm-metaspace.size: 4g

And the pipeline is deployed with these parameters:

--runner=FlinkRunner
--flinkMaster=localhost:8081  
--checkpointingInterval=30000
--parallelism=20
--numConcurrentCheckpoints=200 
--autoBalanceWriteFilesShardingEnabled=true

I have a lot of undelivered messages and I can't find a way to solve it.

0 Answers
Related