【发布时间】: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