【问题标题】:Join Kinesis Streams加入 Ki​​nesis Streams
【发布时间】:2018-09-07 13:38:27
【问题描述】:

我有两个 Kinesis 流,我想创建第三个流,它是这两个流的交集。我的目标是让流处理器响应生成的第三个流上的事件,而无需编写执行此交集的使用者。

stream a 上的记录将是:

{
    "customer_id": 3,
    "first_name":"Marcy",
    "last_name":"Shurtleff"
}

stream b 上的记录是:

{
    "payment_id": 10001,
    "customer_id": 1,
    "amount":234.56,
    "date":"2018-09-07T10:25:43.511Z"

}

我想执行一个连接(就像我可以在带有 Kafka 的 KSQL 中那样)将流 a.customer_id 连接到流 b.customer_id 导致:

{
    "customer_id": 3,
    "first_name":"Marcy",
    "last_name":"Shurtleff",
    "payment_id": 10001,
    "amount":234.56,
    "date":"2018-09-07T10:25:43.511Z"
}

(或我选择的任何类似 sql 的投影)。

我知道 Kafka 和 KSQL 可以做到这一点,但 Kinesis 可以做到吗?

Kinesis Data Analytics 无济于事,因为您不能在该产品中使用多个流作为数据源,并且只能对“应用程序内”流执行联接。

【问题讨论】:

  • Spark 和 Drools 也可以,但不幸的是 Kinesys Analytics 不行

标签: amazon-web-services join amazon-kinesis


【解决方案1】:

我最近使用 Kinesis Data Anlytics 实施了一个完全符合您要求的解决方案。实际上,KDA In-application 仅将一个流作为输入数据源;因此,当您处理多组流时,这种限制使得流入 KDA 的数据的模式标准化成为必要。为了解决这些问题,可以在 lambda 内部使用 python sn-p 代码,通过将其整个有效负载转换为 JSON 编码的字符串来展平和标准化任何事件。下图显示了我的整个解决方案是如何部署的:

流的标准化和扁平化过程详细说明如下:

请注意,在此阶段之后,两个 JSON 事件具有相同的架构且没有嵌套字段。然而,所有信息都被保留了下来。此外,ssn 字段放置在标题上,用作 KDA 应用程序内部的连接键。

有关此解决方案的更多信息,请查看我写的这篇文章:https://medium.com/@guilhermeepassos/joining-and-enriching-multiple-sets-of-streaming-data-with-kinesis-data-analytics-24b4088b5846

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2019-12-15
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-01-24
    • 2022-08-24
    • 1970-01-01
    相关资源
    最近更新 更多