【问题标题】:Having consumer issues when RocksDB in flinkRocksDB 在 flink 中出现消费者问题
【发布时间】: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


    【解决方案1】:

    state.backend.rocksdb.localdir 应该是绝对路径,而不是相对路径。而且这个设置不是用来指定检查点去哪里(不应该在本地磁盘上),这个设置是用来指定工作状态保存在哪里(应该在本地磁盘上)。

    您的工作正面临背压,这意味着管道的某些部分无法跟上。造成背压的最常见原因是 (1) 接收器跟不上,以及 (2) 资源不足(例如,并行度太低)。

    您可以通过使用丢弃接收器运行作业来测试 postgres 是否存在问题。

    查看各种指标应该可以让您了解哪些资源可能配置不足。

    【讨论】:

    • 我在 Flink Dashboard 中检查背压,但没有一个操作员有背压,至少指标是这么说的,但不确定这是否 100% 正确。我只有 4 个 CPU 内核,你认为增加并行度是个好主意吗?这个state.backend.rocksdb.localdir 现在被禁用了。我可以通过丢弃 Sink 来运行测试,但不确定这是否是解决方案,但我会的。
    • rocksdb localdir 默认为 /tmp。
    • 你认为增加并行度是个好主意吗?
    • Flink Dashboard 中的背压监测器并不是一个完全可靠的指标;它可能会说事情还可以,但实际上不是。但是您是否检查了每个子任务中的背压?如果 Flink 源跟不上来自 RabbitMQ 的输入,那么就会出现背压。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2018-10-15
    • 1970-01-01
    • 1970-01-01
    • 2018-07-09
    • 2018-12-31
    • 2017-08-18
    • 1970-01-01
    相关资源
    最近更新 更多