【问题标题】:Kafka Key access on Ingress of a Python Flink Stateful functionPython Flink有状态函数的Ingress上的Kafka Key访问
【发布时间】:2021-12-29 04:09:20
【问题描述】:
我一直在研究 Flink Stateful Functions。它看起来很有希望——除了一件事——我希望我只是想念它。
在我的一生中,无法从 Python 中的 kafka 入口访问 kafka 密钥。在 Java 中,我看到我可以使用反序列化器并将其有效地打包到解码的 message 对象中。但我找不到替代方案。
在我们的例子中,键包含有价值的信息,但值中不存在。
有人遇到过这个 - 还是我错过了?
【问题讨论】:
标签:
python
apache-flink
remote-execution
flink-statefun
【解决方案1】:
首先我要提到的是,用于远程功能的内置 Kafka 入口要求密钥可以解释为 UTF-8 字符串。
如果这确实是你的情况,那么你可以简单地通过以下方式获得它:
def example(context, message):
key = context.address.id
...
print(key) # utf8 string
如果不是这样,那么不幸的是,此时您将不得不使用嵌入式 SDK,在将消息转发到远程函数之前提取键和值。