【问题标题】:How can I access states computed from an external Flink job ? (without knowing its id)如何访问从外部 Flink 作业计算的状态? (不知道它的id)
【发布时间】:2018-02-03 03:19:33
【问题描述】:

我是 Flink 的新手,我目前正在测试一个用例的框架,该用例包括丰富来自 Kafka 的事务,具有许多历史特征(例如,同一源和同一目标之间过去的事务数),然后为这个评分使用机器学习模型进行交易。

目前,功能都保持在 Flink 状态中,并且相同的工作是对丰富的事务进行评分。但我想将特征计算工作与评分工作分开,我不知道该怎么做。

  1. 可查询状态似乎不适合这个,因为需要作业 ID,但如果我错了,请告诉我!

  2. 我曾想过直接查询 RocksDB,但也许有更简单的方法?

  3. 对于 Flink 来说,将这项任务分成两个工作是一个坏主意吗?我们这样做是为了与 Kafka Streams 进行相同的测试,以避免复杂的工作(并检查它是否对延迟有任何积极影响)

一些额外信息:我正在使用 Flink 1.3(但如果需要,我愿意升级)并且代码是用 Scala 编写的

提前感谢您的帮助!

【问题讨论】:

    标签: apache-flink


    【解决方案1】:

    像 Kafka 这样的东西非常适合这种解耦。通过这种方式,您可以拥有一项计算特征并将它们流式传输到 Kafka 主题的工作,该主题由进行评分的工作使用。 (顺便说一句:这样可以很容易地运行几个不同的模型并比较它们的结果。)

    有时使用的另一种方法是调用外部 API 进行评分。 Async I/O 在这里可能会有所帮助。至少有几个小组正在使用流 SQL 来计算特征,并将外部模型评分服务包装为 UDF。

    如果你确实想使用可查询状态,你可以use Flink's REST api to determine the job id

    在 Flink Forward 会议上已经有几次关于在 Flink 中使用机器学习模型的演讲。一个例子:Fast Data at ING – Building a Streaming Data Platform with Flink and Kafka

    社区正在努力使这一切变得更容易。详情请见FLIP-23 - Model Serving

    【讨论】:

    • 感谢您的回答。 Kafka 解决方案看起来不错,但我担心序列化/反序列化(+ 加密)步骤会增加延迟,你怎么看?
    • 现在我想我会调查 Rocksdb 的异步查询(因为最耗时的部分是特征计算,没有太多模型评分)另外,我查看了有关 ING 用例的博客文章,但我觉得他们只做一份工作
    • ING 使用两阶段管道,Kafka 介于两者之间。第一个 flink 作业做数据清洗;第二个计算特征并对模型进行评分。他们在这次谈话中谈到了这些细节:sf-2017.flink-forward.org/kb_sessions/…
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2014-11-21
    • 2020-12-20
    • 1970-01-01
    • 2018-11-02
    • 2022-11-01
    • 2021-09-07
    • 2019-02-27
    相关资源
    最近更新 更多