【发布时间】:2020-09-10 13:35:51
【问题描述】:
我有一份使用 RabbitMQ 的工作,我使用的是 FS 状态后端,但状态的大小似乎变大了,然后我决定将我的状态移动到 RocksDB。 问题是,在运行作业的最初几个小时内,如果流量变慢,则在更多时间后发生事件,但是当流量再次变高时,消费者开始出现问题(事件被堆积为未确认),然后这些问题是反映在应用程序的其余部分。
我有:
4个CPU核心
本地磁盘
16GB 内存
Unix环境
Flink 1.11
Scala 版本 2.11
1 个使用少量 keyedStreams 运行的单个作业,以及大约 10 次转换,然后沉入 Postgres
一些配置
flink.buffer_timeout=50
flink.maxparallelism=4
flink.memory=16
flink.cpu.cores=4
#checkpoints
flink.checkpointing_compression=true
flink.checkpointing_min_pause=30000
flink.checkpointing_timeout=120000
flink.checkpointing_enabled=true
flink.checkpointing_time=60000
flink.max_current_checkpoint=1
#RocksDB configuration
state.backend.rocksdb.localdir=home/username/checkpoints (this is not working don't know why)
state.backend.rocksdb.thread.numfactory=4
state.backend.rocksdb.block.blocksize=16kb
state.backend.rocksdb.block.cache-size=512mb
#rocksdb or heap
state.backend.rocksdb.timer-service.factory=heap (I have test with rocksdb too and is the same)
state.backend.rocksdb.predefined-options=SPINNING_DISK_OPTIMIZED
如果需要更多信息,请告诉我?
【问题讨论】:
标签: apache-flink flink-streaming rocksdb flink-cep rocksdb-java