【问题标题】:Accessing Kafka Topic with two process使用两个进程访问 Kafka 主题
【发布时间】:2019-06-26 16:37:23
【问题描述】:

我有一个工作正常的 Kafka 生产者类。生产者填充 Kafka 主题。其代码如下:

public class kafka_test {
private final static String TOPIC = "flinkTopic";
private final static String BOOTSTRAP_SERVERS = "10.32.0.2:9092,10.32.0.3:9092,10.32.0.4:9092";
public FlinkKafkaConsumer<String> createStringConsumerForTopic(
        String topic, String kafkaAddress, String kafkaGroup) {
    //        ************************** KAFKA Properties ******
    Properties props = new Properties();
    props.setProperty("bootstrap.servers", kafkaAddress);
    props.setProperty("group.id", kafkaGroup);
    FlinkKafkaConsumer<String> myconsumer = new FlinkKafkaConsumer<>(
            topic, new SimpleStringSchema(), props);
    myconsumer.setStartFromLatest();
    return myconsumer;
}
  private static Producer<Long, String> createProducer() {
    Properties props = new Properties();
    props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP_SERVERS);
    props.put(ProducerConfig.CLIENT_ID_CONFIG, "MyKafkaProducer");
    props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, LongSerializer.class.getName());
    props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
    return new KafkaProducer<>(props);
}

public void runProducer(String msg) throws Exception {
    final Producer<Long, String> producer = createProducer();

    try {
            final ProducerRecord<Long, String> record = new ProducerRecord<>(TOPIC, msg );
            RecordMetadata metadata = producer.send(record).get();
            System.out.printf("sent record(key=%s value='%s')" + " metadata(partition=%d, offset=%d)\n",
                    record.key(), record.value(), metadata.partition(), metadata.offset());
    } finally {
        producer.flush();
        producer.close();
    }
 }
}

  public class producerTest {
  public static void main(String[] args) throws Exception{
    kafka_test objKafka=new kafka_test();
    String pathFile="/home/cfms11/IdeaProjects/pooyaflink2/KafkaTest/quickstart/lastDay4.csv";
    String delimiter="\n";
   objKafka.createStringProducer("flinkTopic",
   "10.32.0.2:9092,10.32.0.3:9092,10.32.0.4:9092");
    Scanner scanner = new Scanner(new File(pathFile));
    scanner.useDelimiter(delimiter);
    int i=0;
    while(scanner.hasNext()){
        if (i==0)
            TimeUnit.MINUTES.sleep(1);
         objKafka.runProducer(scanner.next());
       i++;
    }
    scanner.close();
    }
   }

因为我想为我的 Flink 程序提供数据,所以我使用了 Kafka。事实上,我有这部分代码来消费来自 Kafka 主题的数据:

    Properties props = new Properties();
    props.setProperty("bootstrap.servers", 
    "10.32.0.2:9092,10.32.0.3:9092,10.32.0.4:9092");
    props.setProperty("group.id", kafkaGroup);
    FlinkKafkaConsumer<String> myconsumer = new FlinkKafkaConsumer<>(
            "flinkTopic", new SimpleStringSchema(), props);
    DataStream<String> text =   env.addSource(myconsumer).setStartFromEarliest());

我想在我的程序运行的同时运行生产者代码。我的目标是生产者向主题发送一条记录,消费者可以同时从主题轮询该记录。

请您告诉我这是怎么可能的以及如何管理它。

【问题讨论】:

    标签: java apache-kafka apache-flink


    【解决方案1】:

    我认为你需要创建两个类文件,一个是生产者,另一个是消费者。先创建topic再运行consumer,或者直接运行producer。

    【讨论】:

      猜你喜欢
      • 2020-01-05
      • 2019-06-17
      • 2019-07-18
      • 2018-10-27
      • 2020-01-07
      • 1970-01-01
      • 2020-04-12
      • 2015-05-03
      • 2011-04-09
      相关资源
      最近更新 更多