【发布时间】:2020-05-13 01:53:07
【问题描述】:
我需要建议我的项目使用哪些正确的 Kafka 结构以及原因。
我的项目 我正在为投资机器人管理创建一个平台。非常高级 - 您可以编写多个投资策略,并将它们上传到平台,它们将实时执行,提供分析和实时绩效信息。这些策略从 4 个数据流中获取信息。当策略从 4 个不同的 Kafka 主题中读取数据时,这些数据会被传递给策略。这个 kafka 主题直接从交易所 websocket 接收信息。在任何给定时间,平台中都有动态数量的机器人。
我所做的如下: 使用镜像 Kafka-wurmeister 和 zookeper 初始化 kafka 预先初始化我需要的所有 Kakfka 主题。 我通过以下方式将所有信息生成到主题中,将所需数据推送到 Kafka:
payloads = [
{ topic: topic, messages: JSON.stringify(message), partition: 0 }
]
await producer.send(payloads, async function (err, data) {
})
然后我通过一个简单的消费者从主题中读取策略,如下所示: 消费者=新消费者(客户端,[{主题:主题,分区:0}]); consumer.on('message', function (message) {
// Parse the value consumed from kafka
parsedPrice = JSON.parse(message.value)
})
目的是讨论如何使用 kafka 来确保我可以,首先访问来自多个不同消费者的主题,其次有足够的冗余以确保我有非常长的正常运行时间。
【问题讨论】:
标签: apache-kafka kafka-consumer-api kafka-producer-api