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.