我们正在运行一个数据流工作流,它使用来自Kafka的数据,并使用apache写API将时髦的avro文件写入gcs。我们已经配置了最多13个工作者,这些工作者应该处理50k qps的传入事件。我们使用LogAppendTime来处理kafka消息。每条记录的大小彼此相似。窗口为1小时,触发器如下:
Repeatedly
.forever(
AfterFirst.of(
AfterPane.elementCountAtLeast(50000),
AfterProcessingTime.pastFirstElementInPane().plusDelayOf(Duration.standardMinutes(5))
)
)
.orFinally(AfterWatermark.pastEndOfWindow())事件正在以30k qps的速率生成到Kafka。过了一段时间,工作人员只能处理25k qps,数据水印延迟增加,延迟约2小时。因此,我们更新了工作流,并假设它可以解决问题。更新后,在最初的一个小时内,它能够像预期的那样处理45k qps,并且数据水印正在减少。超过这一点,大多数工作进程的CPU利用率下降,qps下降到20k。这导致数据水印延迟随着事件以30k qps速率产生到Kafka而增加。
在进一步的调查中,我们发现大多数工人的CPU利用率都很低,从15%到20%不等,其中2人的CPU利用率为40%,其中一人的CPU利用率为60%。从日志中我们可以看到,CPU使用率较高的CPU使用率较高的CPU使用率较低的gcs写入gcs的频率更高。我们将numShards设置为26,假设碎片将均匀分布在工作进程中。但是,看起来数据流将它们中的大多数分配给相同的工作进程。
工作流的详细信息:
job_id: 2018-05-08_10_39_57-9264166384462032078
numShards: 26
maxNumWorkers: 13
发布于 2018-05-19 15:50:28
只是不想让这个问题悬而未决。正如@revathy评论的那样,这个问题是通过从标准永久磁盘更改为SSD永久磁盘来修复的,因为工人受到永久磁盘可以执行的IOPs数量的限制。
https://stackoverflow.com/questions/50242762
复制相似问题