【问题标题】:How to configure Apache Kafka to sending data at specified time?如何配置 Apache Kafka 在指定时间发送数据?
【发布时间】:2017-12-10 06:31:56
【问题描述】:

让我们考虑以下工作 Apache Kafka 的抽象架构:

Poducers ->(Send messages) -> Apache Kafka -> (Resend to customers) -> Customers 

可以配置Kafka在指定时间给客户发送消息吗?

第二个问题,客户回滚消息到Kafka是真的吗?

【问题讨论】:

  • 你能解释一下“回滚消息到 Kafka”是什么意思吗?你的意思是阅读之前的消息?还是将来自消费者的消息传回 Kafka?

标签: apache-kafka kafka-consumer-api


【解决方案1】:

如果我收到问题,您想在特定时间点向客户发送数据。如果在 Apache Kafka 中使用 Lenses,它可能很简单

#cron the following to execute daily at 24:00 curl -XGET http://lenses-host:port/api/sql/data?sql=SELECT * from topicA WHERE customer = 'customerA WHERE _ts > 'yyyy-mm-dd hh:mm:ss'' > customerA.json send info@customerA.com customerA.json

因此,要回答问题的第一部分,您需要构建消费者逻辑。 Kafka 不支持回滚,但您可以轻松地执行以下操作:

INSERT INTO topicB SELECT * from topicA WHERE _ts < '2017-12-10 00:00:00'

因此您可以轻松地从另一个主题创建新主题,但没有回滚语义。

【讨论】:

    【解决方案2】:

    消费者从 Kafka 拉取消息; Kafka 不会推送(“发送”)它们。因此,您的消费者可以在需要时提取数据。

    【讨论】:

    • 那么,如何每次(每毫秒)拉取数据并将消息回滚到Kafka?客户可以按条件提取数据吗?
    【解决方案3】:

    正如其他人已经回答的那样,Kafka 不会将消息推送给消费者,而是消费者从 Kafka 中提取消息;这意味着您需要编写消费者以便在特定时间(或间隔)从 Kafka 主题中提取消息。 关于回滚是什么意思?也许该消费者从 Kafka 获取消息,但随后又想重新读取相同的消息,因为在第一次处理期间发生错误?如果是,关于 Kafka 有两个方面需要考虑:

    • Kafka 保留可配置的消息(甚至数天),这意味着当消费者收到消息时,它们不会从主题分区中删除
    • 相反,当消费者收到消息时,它必须提交偏移量,以便它可以跟踪从主题分区读取的最新消息是什么。此提交可以自动或手动完成,以便您只有在过程顺利时才能提交偏移量。在任何情况下,您都可以重绕流并决定从特定偏移量重新开始读取主题分区。

    【讨论】:

    • 其实你就在这里; it means that you need to write your consumer in order to pull messages 。这是一个原始的问题,这使我想到了卡夫卡。你能为发布者提供另一种机制吗?
    【解决方案4】:

    延伸到罗宾,

    Kafka 不会向消费者推送消息,消费者需要从 Kafka 拉取消息。

    看下面python sn-p从Kafka读取消息:

    running = True
    while running:
        msg = c.poll(timeout=1.0)
        if not msg.error():
            print('Received message: %s' % msg.value().decode('utf-8'))
        elif msg.error().code() != KafkaError._PARTITION_EOF:
            print(msg.error())
            running = False
    

    上面的 sn-p msg = c.poll(timeout=1.0) 用于每秒从 Kafka 拉消息。如果您想将超时增加到任意秒数。这意味着 Kafka 消费者消费者将从每个时间间隔提取消息。

    如果你想做调度,你必须在调度时间调用 poll 方法。

    注意:你的 session.timeout.ms 应该大于 poll time

    【讨论】:

    • 我不问如何在客户端阅读消息。我问了另一个,请再读一遍问题
    • 您只回答了部分问题,如果在 Kafka 中可以回滚,您什么也没说。
    • 回滚是什么意思??你到底是什么意思?你能举个例子吗?客户意味着不是消费者?
    猜你喜欢
    • 1970-01-01
    • 2019-01-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-09-18
    相关资源
    最近更新 更多