【问题标题】:is this the correct way to read the message through kafka producer and push it to the topic这是通过kafka生产者读取消息并将其推送到主题的正确方法吗
【发布时间】:2018-11-08 16:09:34
【问题描述】:

我已经编写了这个 Kafka 生产者并从桌面读取一个文件,然后将文件中的数据作为值推送,并通过每次读取每一行时添加一个来生成我自己的密钥。这是正确的方法还是我做了我不应该做的事情?请需要一些建议。 我可以在我的主题中看到消息,但每个消息都与一个键相关,所以如果我有一个用例,如果我从外部读取它,我可以推送任何这样的日志数据。我可以使用日志数据作为值,还是应该采用完全不同的逻辑。 请帮忙

import java.io.BufferedReader;
import java.io.File;
import java.io.FileReader;
import java.io.IOException;
import java.util.Properties;

import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.serialization.StringSerializer;


public class SyncProducer {

    public static void main(String[] args) throws IOException {

        File file = new File("/Users/adityaverma/Desktop/ParseData.txt");

        BufferedReader br = new BufferedReader(new FileReader(file));

        Properties properties  = new Properties();
        properties.setProperty("bootstrap.servers","127.0.0.1:9092");
        properties.setProperty("key.serializer",StringSerializer.class.getName()); // our key and values are String
        properties.setProperty("value.serializer",StringSerializer.class.getName());
        properties.setProperty("acks", "1"); 
        properties.setProperty("retries", "3"); 
        properties.setProperty("linger.ms", "1"); 



        Producer<String,String> producer = new org.apache.kafka.clients.producer.KafkaProducer<String,String>(properties);
        // these will go in random partition as we increment the key

        String line = " ";

        int key = 0;
        while((line = br.readLine()) != null){
        //  System.out.println(line);

        ProducerRecord<String,String> producerRecord = new ProducerRecord<String,String>("try_Buffered3Part",Integer.toString(key),line);
         key++;
         System.out.println(key);
        producer.send(producerRecord);

        }
        producer.close();
        System.out.println("exit");
    }

}

【问题讨论】:

    标签: java file apache-kafka hadoop-yarn hadoop2


    【解决方案1】:

    看起来不错。如果你对随机分区没问题,你可以选择使用空键。

    您可能还想查看 logstash 和 kafka 的集成。

    https://www.elastic.co/guide/en/logstash/current/plugins-outputs-kafka.html

    【讨论】:

    • 我明白了,谢谢,是的,我知道随机密钥。感谢分享链接。所以我的下一个问题是,如果我将它们放入 HDFS 将如何读取?像我可以根据密钥为每个分区创建一个目录吗?然后我可能只想要将其放入蜂巢表的值。另外,我知道有一个 Kafka 连接可以将数据放入 HDFS,但我将如何以编程方式进行操作。现在我正在努力通过eclipse将数据推送到hdfs而不创建jar文件然后将该jar文件部署到Hadoop集群中然后从那里运行它。我想从eclipse做
    • 你想把日志、kafka、hdfs 还是 hive 放在哪里?如果不是kafka,要么更新问题,要么创建一个新问题。
    • 所以我换个说法,我想通过Kafka从某个来源读取日志,然后将其写入HDFS。
    • 您也可以直接将日志文件复制到hdfs。它会更有效率,尤其是在文件轮换之后。
    • 如果生成了大量日志文件并且我需要将数据实时摄取到 HDFS 中怎么办。我在这里担心的是,kafka 在从多个分区读取方面会更好地扩展
    【解决方案2】:

    是不是正确的方式

    不清楚你的目标。你可以从终端消费数据吗?那你生产得很好。

    您可以使用整数作为键。 Kafka 有一个 IntegerSerializer

    使用 null 作为键或排除该参数是向随机分区发送数据的标准方式,并且不会遇到整数过载

    我想通过Kafka从某个来源读取日志,然后将其写入HDFS。

    如果您只想将数据记录到 Hadoop 中,Fluentd 或 Logstash 可以做到这一点。

    在开始使用 Kafka 走这条路之前,您绝对应该选择一种数据格式。例如,Hadoop 和 Kafka 更喜欢 Avro 或 JSON 而不是 CSV。 Confluent 有大量关于将 Avro 生成到 Kafka 中的文档

    您可以使用 Kafka Connect HDFS 连接器或 Apache Nifi 将 Kafka 数据导入 Hadoop。不要重新发明轮子来编写自己的消费者。

    【讨论】:

    • 我明白了,我实际上是 Kafka 环境的新手,并且正在尝试逐步学习。谢谢你让我知道,我不知道数据格式的概念,谢谢你为那部分带来了光明,我很感激。对于 Nifi,我应该知道什么以及如何掌握?
    • 是的,我可以从终端消费数据
    • Nifi 只是另一个开源软件。但它有一个 GUI 供您拖放流式处理...有些人更喜欢编写 Kafka 代码,如果您使用的是 Hortonworks HDF Hadoop 集群,您可以点击几下安装它...如果您只想留在 Kafka 生态系统中,然后查看docs.confluent.io/current/connect/connect-hdfs/docs/index.html
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-01-17
    • 2020-08-19
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多