0Pricing
Apache Kafka & Stream Processing Fundamentals · 课时

向 Kafka 生产消息

探索如何编写应用,高效且可靠地向 Kafka 主题发送数据

向 Kafka 生产消息 是 CoddyKit 上的免费 Apache Kafka & Stream Processing Fundamentals 课时。 这是第 1 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 Apache Kafka & Stream Processing Fundamentals 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 Apache Kafka & Stream Processing Fundamentals 课程共包含 4 节课。

本课时的部分内容尚未翻译,以英文显示。

Meet the Kafka Producer

In this lesson, we'll learn how to send messages (also called records) to Kafka topics. This is the job of a Kafka Producer.

Think of a producer as an application or service that generates data. It then sends this data to a Kafka cluster, where it's stored in specific topics for other applications to read.

  • Producers generate data.
  • Topics organize data streams.
  • Brokers store the data.

Producer Client Basics

To send data, your application uses a Kafka Producer client library. This library handles all the complex interactions with the Kafka brokers.

It takes care of:

  • Finding the right Kafka broker.
  • Serializing your data into bytes.
  • Retrying failed send operations.
  • Balancing message distribution.

We'll use Java examples, but the concepts apply across languages.

Essential Producer Config

Before a producer can send messages, it needs some basic configuration. The two most important settings are:

  • bootstrap.servers: A comma-separated list of host/port pairs for Kafka brokers. The producer uses these to discover the full cluster.
  • key.serializer and value.serializer: Classes that convert your message's key and value objects into byte arrays, which is how Kafka stores data.

Without these, the producer won't know where to send messages or how to format them.

Producer Config Example

Here's how you might set up these properties in Java:

import java.util.Properties;

public class ProducerConfigDemo {
  public static void main(String[] args) {
    Properties props = new Properties();
    props.put("bootstrap.servers", "localhost:9092");
    props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
    props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");

    System.out.println("Producer properties configured!");
    // In a real app, you'd create a KafkaProducer with these props
  }
}

Crafting Your Message: ProducerRecord

When you send data to Kafka, you don't just send a string. You send a ProducerRecord. This object encapsulates your message and its metadata.

A ProducerRecord requires:

  • The topic name where the message will be sent.
  • An optional key: Used for partitioning messages. Messages with the same key go to the same partition.
  • The value: The actual data you want to send.

Keys are important for ensuring ordering for related data within a topic.

Sending Messages (Blocking)

The simplest way to send a message is using the send() method. If you want to wait for the message to be acknowledged by Kafka, you can call .get() on the returned Future object.

This makes the send operation synchronous (blocking). It's easy to understand, but can be slow if you're sending many messages.

import org.apache.kafka.clients.producer.*;
import java.util.Properties;

public class SyncProducer {
  public static void main(String[] args) throws Exception {
    Properties props = new Properties();
    props.put("bootstrap.servers", "localhost:9092");
    props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
    props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");

    try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) {
      ProducerRecord<String, String> record = new ProducerRecord<>(
        "my-topic", "key1", "Hello Sync Kafka!");
      
      RecordMetadata metadata = producer.send(record).get(); // Blocks here
      System.out.println("Sent message: " + metadata.topic() + "-" + metadata.partition());
    }
  }
}

Sending Messages (Non-Blocking)

For better performance, Kafka producers are designed to send messages asynchronously. When you call send(), it adds the message to a buffer and returns immediately.

The actual sending happens in the background. This allows your application to continue processing without waiting for each message to be delivered.

import org.apache.kafka.clients.producer.*;
import java.util.Properties;

public class AsyncProducer {
  public static void main(String[] args) {
    Properties props = new Properties();
    props.put("bootstrap.servers", "localhost:9092");
    props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
    props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");

    try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) {
      ProducerRecord<String, String> record = new ProducerRecord<>(
        "my-topic", "key2", "Hello Async Kafka!");
      
      producer.send(record); // Returns immediately
      System.out.println("Message queued for sending.");
      // In a real app, you'd send many messages here
    }
  }
}

Handling Send Results with Callbacks

Since send() is asynchronous, how do you know if a message was successfully sent or if an error occurred? You use a Callback.

The callback function is executed once Kafka acknowledges the message or if an error prevents it from being sent. This is crucial for error handling and logging.

import org.apache.kafka.clients.producer.*;
import java.util.Properties;

public class CallbackProducer {
  public static void main(String[] args) {
    Properties props = new Properties();
    props.put("bootstrap.servers", "localhost:9092");
    props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
    props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");

    try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) {
      ProducerRecord<String, String> record = new ProducerRecord<>(
        "my-topic", "key3", "Hello Callback Kafka!");
      
      producer.send(record, new Callback() {
        @Override
        public void onCompletion(RecordMetadata metadata, Exception exception) {
          if (exception == null) {
            System.out.println("Message sent successfully to topic " + metadata.topic());
          } else {
            System.err.println("Error sending message: " + exception.getMessage());
          }
        }
      });
      // Must flush or close producer to ensure callback is triggered in short programs
      producer.flush(); 
    }
  }
}

Ensuring Delivery: Acks Configuration

Producer reliability is controlled by the acks configuration. This setting determines how many acknowledgments a producer needs from Kafka brokers before considering a message 'sent'.

  • acks=0: Producer sends and doesn't wait for any acknowledgment. Fastest, but lowest durability (messages might be lost).
  • acks=1: Producer waits for the leader broker to acknowledge receipt. Good balance of speed and durability.
  • acks=all (or -1): Producer waits for all in-sync replicas to acknowledge. Slowest, but highest durability (messages are very unlikely to be lost).

Producer Best Practices

To ensure your Kafka producers are efficient and robust:

  • Always close the producer: Call producer.close() when your application shuts down. This flushes any buffered messages and releases resources.
  • Batching: Kafka producers automatically batch messages for efficiency. You can tune linger.ms and batch.size for optimal throughput.
  • Error Handling: Implement robust error handling in your callbacks to deal with transient network issues or permanent errors.

Proper configuration and resource management are key to a healthy Kafka application.

Quick Check: Producers

Which producer configuration property specifies how many acknowledgments the producer needs from Kafka brokers before considering a message successfully sent?

Producer Summary

You've learned the fundamentals of producing messages to Kafka!

  • Producers send data to topics.
  • Key configurations include bootstrap.servers and serializers.
  • Messages are wrapped in ProducerRecord objects.
  • You can send messages synchronously (blocking) or asynchronously (non-blocking).
  • Callbacks are used to handle asynchronous send results.
  • The acks setting controls message durability.

Next, we'll dive into how applications read these messages using Kafka Consumers!

常见问题解答

「向 Kafka 生产消息」课时是免费的吗?

是的 — 「向 Kafka 生产消息」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 Apache Kafka & Stream Processing Fundamentals 课程的其余内容,请升级到 CoddyKit PRO。 Apache Kafka & Stream Processing Fundamentals 课程共包含 4 节课。

「向 Kafka 生产消息」这节课中我会学到什么?

探索如何编写应用,高效且可靠地向 Kafka 主题发送数据 你通过在浏览器中直接运行的动手代码来练习 Apache Kafka & Stream Processing Fundamentals,全天候 AI 导师会在你学习这节课的过程中回答你的问题。

学习 Apache Kafka & Stream Processing Fundamentals 需要有经验吗?

无需任何先前经验。CoddyKit 上的 Apache Kafka & Stream Processing Fundamentals 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 1 节课,共 4 节。

「向 Kafka 生产消息」课时需要多长时间?

大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。

我能在这节 Apache Kafka & Stream Processing Fundamentals 课中编写并运行代码吗?

能。每节 Apache Kafka & Stream Processing Fundamentals 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。

此课程中的所有课时

  1. 向 Kafka 生产消息
  2. 从 Kafka 消费消息
  3. 理解分区和偏移量
  4. 消息键与分区策略
← 返回 Apache Kafka & Stream Processing Fundamentals