【问题标题】:Accessing kafka timestamps in pyflink在pyflink中访问kafka时间戳
【发布时间】:2021-07-16 05:13:57
【问题描述】:

我正在尝试编写一个 Pyflink 应用程序来测量延迟和吞吐量。我的数据来自 kafka 主题的 json 对象,并使用 SimpleStringSchema-class 加载到 DataStream 中以进行反序列化。按照这篇文章的答案 (How performance can be tested in Kafka and Flink environment?),我让 Kafka 生产者在事件中添加了时间戳,但现在很难理解如何访问这些时间戳。我知道上面提到的帖子为这个问题提供了一个解决方案,但我正在努力将此示例转移到 python,因为文档/示例很少。

另一个帖子 (Apache Flink: How to get timestamp of events in ingestion time mode?) 建议我应该定义一个 ProcessFunction。但是,在这里我也不确定语法。我可能不得不做这样的事情(取自:https://github.com/apache/flink/blob/master/flink-end-to-end-tests/flink-python-test/python/datastream/data_stream_job.py)

class MyProcessFunction():

    def process_element(self, value, ctx):
        result = value.get_time_stamp()
        yield result

在这里value.get_time_stamp() 的正确方法是什么?或者有没有更简单的方法可以解决我不知道的问题?

谢谢!

【问题讨论】:

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


    【解决方案1】:

    当您设置一个由 Kafka 主题支持的表时,您可以为 Kafka 时间戳声明一个虚拟列,如本示例中的 event_time 列:

    CREATE TABLE KafkaTable (
      `event_time` TIMESTAMP(3) METADATA FROM 'timestamp',
      `partition` BIGINT METADATA VIRTUAL,
      `offset` BIGINT METADATA VIRTUAL,
      `user_id` BIGINT,
      `item_id` BIGINT,
      `behavior` STRING
    ) WITH (
      'connector' = 'kafka',
      'topic' = 'user_behavior',
      'properties.bootstrap.servers' = 'localhost:9092',
      'properties.group.id' = 'testGroup',
      'scan.startup.mode' = 'earliest-offset',
      'format' = 'csv'
    );
    

    有关使用 Kafka 标头中的元数据的更多信息,请参阅 documentation for Flink's Kafka Table connector。

    【讨论】:

    • 非常感谢您的回答。有没有办法用 DataStream-API 做到这一点?
    • 我不相信,不。
    猜你喜欢
    • 2021-11-23
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-02-16
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多