【问题标题】:Picks on latency graph of a Flink/Kafka application挑选 Flink/Kafka 应用程序的延迟图
【发布时间】:2020-01-27 12:15:29
【问题描述】:

我有一个应用程序从 Kafka 主题接收推文,有一个一秒钟的窗口,然后通过 AsyncIO 操作将这些推文保存在 Cassandra 上,该操作允许打开最多 100 个线程(AsyncIO 操作符的最后一个参数)而不做对数据进行任何预处理:只需逐条保存推文,并附上保存时间的时间戳。

然后,我强调了 Flink 应用程序发送了 300 万条推文,并在 Grafana 中做了一个图表,显示了数据库中保存了多少条推文,但是这个图表显示了一些选择,不是一条连续的线,我可以'不明白为什么。

因此您可以看到,在一分钟的时间间隔内,它会节省 7k,然后再节省 5k,然后再节省 2k。我怎样才能找出原因?

【问题讨论】:

  • 为什么不使用 Kafka Connect 而不是 AsyncIO?
  • 因为我是这样做的,所以我不知道Kafka connect,会搜索更多关于它的信息。另外,以后我会使用flink,以便在保存到数据库之前对数据进行一些操作。您知道为什么会出现这种行为吗?
  • 你使用的是事件时间还是处理时间语义?
  • 我不知道...另外,不清楚您是直播推文还是批量加载。此外,默认情况下,Kafka 生产者会将事件一起批处理,因此您可以使用批处理或请求大小。 RE: Connect,你仍然可以从 Kafka->Kafka 使用 Flink,然后从那里使用 Kafka Connect 进入数据库
  • 直播推文上的事件时间语义。 Kafka producer 上的批处理事件我几个小时前才弄清楚,玩它是个好主意。我会试试的。

标签: twitter apache-kafka apache-flink flink-streaming


【解决方案1】:

首先,如果你想写信给 cassandra,我会使用connector。如果不是几乎不可能的话,手动正确地实现完全一次是非常困难的。

其次,AsyncIO 没有启动 100 个线程。事实上,它并没有为用户启动任何线程。您需要通过任何方式自己启动它们。通常,它使用库有自己的连接池的外部系统的回调机制。

如果您正在进行同步调用,则需要管理自己的线程池。我建议使用 Executors.newCachedThreadPool() 并将您的异步任务提交给它。 AsyncIO 只会帮助将异步结果合并回同步流中。

第三,100 个线程可能很多,具体取决于您的设置。另请注意,如果您使用 Flink 的纵向扩展(每个任务管理器使用多个插槽),则使用的线程会成倍增加。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2019-10-08
    • 2015-01-23
    • 2022-01-09
    • 1970-01-01
    • 2021-10-16
    • 2013-12-31
    • 1970-01-01
    相关资源
    最近更新 更多