Наивный производитель Kafka не работает, когда я пытаюсь создать сообщения, которые читаются из файла

Я пытаюсь написать наивного продюсера Kafka на Java. Приложение принимает два входа:

  1. Название темы Kafka, для которой должны быть созданы сообщения
  2. Путь к файлу, содержащему сообщения, которые будут отправлены в Kafka

Я написал следующий код. Когда я запускаю его, я вижу, что операторы System.out.println печатают ожидаемые значения, но по какой-то причине сообщения не отправляются в Kafka. Что мне нужно изменить, чтобы он заработал?

package com.myname.kafka.producer;

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

import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;

public class NaiveKafkaProducer {

    private static final Properties properties = new Properties();

    private static Producer<String, String> producer;

    private static String topic;

    private static BufferedReader br;

    static {
        properties.put("bootstrap.servers", "localhost:9092");
        properties.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
        properties.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
        properties.put("request.required.acks", "all");
        System.out.println("Creating Kafka producer with the following properties :: " + properties);
        producer = new KafkaProducer<>(properties);
    }

    @SuppressWarnings("resource")
    public static void main(String[] args) throws IOException {
        try {
            if (args.length != 0) {
                topic = args[0];
                File file = new File(args[1]);
                br = new BufferedReader((Reader) new FileReader(file));
            }
        } catch (Exception e) {
            System.out.println("Check input arguments. Error thrown while populating arguments to local variables");
            e.printStackTrace();
        }

        String msg;
        while ((msg = br.readLine()) != null) {
            System.out.println("Message to publish : " + msg);
            System.out.println("Topic : " + topic);
            producer.send(new ProducerRecord<String, String>(topic, "", msg));
        }
        return;
    }
}

На удивление работает следующий код (в котором я все жестко запрограммировал):

package com.myname.kafka.producer;

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

import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;

public class NaiveKafkaProducer {

    private static final Properties properties = new Properties();

    private static Producer<String, String> producer;

    private static String topic;

    private static BufferedReader br;

    static {
        properties.put("bootstrap.servers", "localhost:9092");
        properties.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
        properties.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
        properties.put("request.required.acks", "all");
        System.out.println("Creating Kafka producer with the following properties :: " + properties);
        producer = new KafkaProducer<>(properties);
    }

    public static void main(String[] args) throws IOException {
            try {
                String[] msgs = new String[2];
                msgs[0] = "message 1";
                msgs[1] = "message 2";
                topic = "mytopic"
                for(String msg:msgs){
                    producer.send(new ProducerRecord<String, String>(topic, "", msg));
                }
                producer.close();
            } catch (Exception e) {
                System.out.println("Exception caught in main method while trying to produce the messages to Kafka");
                e.printStackTrace();
            }
    }
}
Пользовательский скаляр GraphQL
Пользовательский скаляр GraphQL
Листовые узлы системы типов GraphQL называются скалярами. Достигнув скалярного типа, невозможно спуститься дальше по иерархии типов. Скалярный тип...
Как вычислять биты и понимать побитовые операторы в Java - объяснение с примерами
Как вычислять биты и понимать побитовые операторы в Java - объяснение с примерами
В компьютерном программировании биты играют важнейшую роль в представлении и манипулировании данными на двоичном уровне. Побитовые операции...
Поднятие тревоги для долго выполняющихся методов в Spring Boot
Поднятие тревоги для долго выполняющихся методов в Spring Boot
Приходилось ли вам сталкиваться с требованиями, в которых вас могли попросить поднять тревогу или выдать ошибку, когда метод Java занимает больше...
Полный курс Java для разработчиков веб-сайтов и приложений
Полный курс Java для разработчиков веб-сайтов и приложений
Получите сертификат Java Web и Application Developer, используя наш курс.
0
0
52
1
Перейти к ответу Данный вопрос помечен как решенный

Ответы 1

Ответ принят как подходящий

критический метод вызван во втором фрагменте и отсутствует в первом

producer.close();

из документации для этого метода:

Close this producer. This method blocks until all previously sent requests complete.

Когда вы вызываете метод produce, на самом деле это не означает, что сообщение было создано. Метод возвращает вам future. Вы можете дождаться создания каждого сообщения, вызывая get() для каждого результата метода создания.

Другие вопросы по теме