【问题标题】:Kafka Consumer which writes into multiple file写入多个文件的Kafka消费者
【发布时间】:2018-07-30 16:59:43
【问题描述】:

我必须实现一个 kafka 消费者,它从主题中读取数据并根据有效负载中存在的帐户 ID(将接近百万)将其写入文件。假设每秒大约有 3K 个事件。可以为每条阅读的消息打开和关闭文件吗? 还是我应该考虑不同的方法?

【问题讨论】:

  • 它会是一个文件吗?还是每个帐户 ID 都有自己唯一的文件?
  • 每个 accountid 都有自己唯一的文件。
  • 好的..那么答案应该对你有用

标签: file-io apache-kafka


【解决方案1】:

我假设如下:

  1. 每个帐户 ID 都是唯一的,并且有自己唯一的文件。
  2. 文件中的数据有一点延迟是可以的,即文件中的数据将接近实时。
  3. 每个事件读取的数据并不大。

解决方案:

  1. Kafka 消费者读取数据并写入数据库,最好是 NoSQL 数据库。
  2. 一个单独的单线程定期读取数据库以查找插入的新记录,并按 accountId 对它们进行分组。
  3. 然后遍历 accoundId 并为每个 accountId 打开文件,立即写入数据,关闭文件并移动到下一个 accountId。

优点:

  1. 您的消费者不会因为文件处理而被阻止,因为这两个操作是分离的。
  2. 即使文件处理失败,数据也始终存在于数据库中以进行重新处理。

【讨论】:

  • 这是我想到的替代解决方案,因为这些数据已经被持久化到 cassandra。想探索一个选项,我可以在此操作期间避免 cassandra 出现峰值。主要目标是为每个帐户创建每日文件。
  • 如果数据已经在 Cassandra 中,那么为什么不从那里读取呢。读取操作并不是真正的尖峰。
【解决方案2】:

如果您的帐户 id 重复,那么最好开窗。您可以通过窗口聚合所有事件,例如 1 分钟,然后您可以按键分组事件并一次处理所有 accountId。

这样,您不必多次打开文件。

【讨论】:

  • 无法保证我会多次看到所有帐户 ID。
【解决方案3】:

不能为每条消息都打开一个文件,您应该缓冲固定数量的消息,然后在每个限制时写入文件。


您可以使用 Confluent 提供的 HDFS Kafka 连接器来管理它。

如果配置了FieldPartitioner 写入给定store.url=file:///tmp 的本地文件系统,例如,这将在您的主题中为每个唯一的accountId 字段创建一个目录。然后flush.size 配置决定了单个文件中最终会包含多少条消息

Hadoop 不需要安装,因为 HDFS 库包含在 Kafka Connect 类路径中,并且它们支持本地文件系统

你会在创建两个属性文件后这样启动它

bin/connect-standalone worker.properties hdfs-local-connect.properties 

【讨论】:

  • 我不认为这对我有用。我有一些逻辑可以检查数据是否要写入文件。忘了提到这一点,每个帐户 id 都有自己的目录,文件名将基于日期。
猜你喜欢
  • 2011-11-08
  • 2019-04-07
  • 2017-01-26
  • 1970-01-01
  • 2018-12-31
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-08-31
相关资源
最近更新 更多