【问题标题】:How to count the number of messages fetched from a Kafka topic in a day?如何统计一天从 Kafka 主题获取的消息数?
【发布时间】:2019-11-25 03:19:52
【问题描述】:

我正在从 Kafka 主题中获取数据并以 Deltalake(parquet) 格式存储它们。 我希望找出特定日期获取的消息数

我的思考过程:我想使用 spark 读取数据以 parquet 格式存储的目录,并对特定日期的“.parquet”文件应用计数。这会返回一个计数,但我不确定这是否正确。

这种方式正确吗?有没有其他方法可以计算特定日期(或持续时间)从 Kafka 主题获取的消息数量?

【问题讨论】:

    标签: apache-spark apache-kafka parquet spark-structured-streaming delta-lake


    【解决方案1】:

    我们从主题消费的消息不仅有键值,还有其他信息,如时间戳

    可用于跟踪消费流。

    时间戳 时间戳由 Broker 或 Producer 根据主题配置更新。如果 Topic 配置的时间戳类型为 CREATE_TIME,则 broker 将使用生产者记录中的时间戳,而如果 Topic 配置为 LOG_APPEND_TIME ,则在附加记录时,时间戳将由 broker 使用 broker 本地时间覆盖。

    1. 因此,如果您要存储任何位置,如果您保留时间戳,则可以很好地跟踪每天或每小时的消息速率。

    2. 您可以通过其他方式使用一些 Kafka 仪表板,例如 Confluent Control Center(许可价格)或 Grafana(免费)或任何其他工具来跟踪消息流。

    3. 在我们的例子中,在消费消息和存储或处理消息的同时,我们还将消息的元详细信息路由到 Elastic Search,我们可以通过 Kibana 将其可视化。

    【讨论】:

      【解决方案2】:

      您可以利用 Delta Lake 提供的“时间旅行”功能。

      你可以这样做

      // define location of delta table
      val deltaPath = "file:///tmp/delta/table"
      
      // travel back in time to the start and end of the day using the option 'timestampAsOf'
      val countStart = spark.read.format("delta").option("timestampAsOf", "2021-04-19 00:00:00").load(deltaPath).count()
      val countEnd = spark.read.format("delta").option("timestampAsOf", "2021-04-19 23:59:59").load(deltaPath).count()
      
      // print out the number of messages stored in Delta Table within one day
      println(countEnd - countStart)
      
      

      请参阅Query an older snapshot of a table (time travel) 上的文档。

      【讨论】:

        【解决方案3】:

        在不计算两个版本之间的行数的情况下检索此信息的另一种方法是使用Delta table history。这样做有几个优点 - 您无需读取整个数据集,您也可以考虑更新和删除,例如,如果您正在执行 MERGE 操作(无法在不同版本上比较 .count ,因为更新是替换实际值,或者删除行)。

        例如,对于仅追加,以下代码将计算所有由正常append 操作写入的插入行(对于其他事情,例如,MERGE/UPDATE/DELETE,我们可能需要查看其他metrics):

        from delta.tables import *
        
        df = DeltaTable.forName(spark, "ml_versioning.airbnb").history()\
          .filter("timestamp > 'begin_of_day' and timestamp < 'end_of_day'")\
          .selectExpr("cast(nvl(element_at(operationMetrics, 'numOutputRows'), '0') as long) as rows")\
          .groupBy().sum()
        

        【讨论】:

          猜你喜欢
          • 1970-01-01
          • 2015-08-01
          • 2020-12-09
          • 2019-02-10
          • 2021-08-07
          • 2018-10-27
          • 2019-06-23
          • 2021-12-12
          • 2020-06-19
          相关资源
          最近更新 更多