{"id":872,"date":"2023-03-09T21:45:17","date_gmt":"2023-03-09T13:45:17","guid":{"rendered":"https:\/\/www.appblog.cn\/?p=872"},"modified":"2023-04-29T16:43:35","modified_gmt":"2023-04-29T08:43:35","slug":"spring-boot-integrate-kafka-to-implement-producer-and-consumer","status":"publish","type":"post","link":"https:\/\/www.appblog.cn\/index.php\/2023\/03\/09\/spring-boot-integrate-kafka-to-implement-producer-and-consumer\/","title":{"rendered":"Spring Boot\u96c6\u6210Kafka\u5b9e\u73b0producer\u548cconsumer"},"content":{"rendered":"<p>\u4ecb\u7ecd\u5982\u4f55\u5728Spring Boot\u9879\u76ee\u4e2d\u96c6\u6210Kafka\u6536\u53d1Message<\/p>\n<h2>pom.xml\u4f9d\u8d56<\/h2>\n<p><!-- more --><\/p>\n<pre><code class=\"language-xml\">&lt;!-- https:\/\/mvnrepository.com\/artifact\/org.springframework.kafka\/spring-kafka --&gt;\n&lt;dependency&gt;\n    &lt;groupId&gt;org.springframework.kafka&lt;\/groupId&gt;\n    &lt;artifactId&gt;spring-kafka&lt;\/artifactId&gt;\n    &lt;version&gt;2.2.4.RELEASE&lt;\/version&gt;\n&lt;\/dependency&gt;<\/code><\/pre>\n<h2>\u914d\u7f6e\u6587\u4ef6<\/h2>\n<pre><code>#=================== kafka ===================\nkafka.consumer.zookeeper.connect=127.0.0.1:2181\nkafka.consumer.servers=127.0.0.1:9092\nkafka.consumer.enable.auto.commit=true\nkafka.consumer.session.timeout=6000\nkafka.consumer.auto.commit.interval=100\nkafka.consumer.auto.offset.reset=latest\nkafka.consumer.topic=test\nkafka.consumer.group.id=test\nkafka.consumer.concurrency=10\n\nkafka.producer.servers=127.0.0.1:9092\nkafka.producer.retries=0\nkafka.producer.batch.size=4096\nkafka.producer.linger=1\nkafka.producer.buffer.memory=40960<\/code><\/pre>\n<h2>Kafka producer<\/h2>\n<p>\uff081\uff09\u901a\u8fc7<code>@Configuration<\/code>\u3001<code>@EnableKafka<\/code>\uff0c\u58f0\u660eConfig\u5e76\u4e14\u6253\u5f00KafkaTemplate\u80fd\u529b\u3002<br \/>\n\uff082\uff09\u901a\u8fc7<code>@Value<\/code>\u6ce8\u5165<code>application.properties<\/code>\u914d\u7f6e\u6587\u4ef6\u4e2d\u7684Kafka\u914d\u7f6e\u3002<br \/>\n\uff083\uff09\u901a\u8fc7<code>@Bean<\/code>\u751f\u6210bean<\/p>\n<pre><code class=\"language-java\">import org.apache.kafka.clients.producer.ProducerConfig;\nimport org.apache.kafka.common.serialization.StringSerializer;\nimport org.springframework.beans.factory.annotation.Value;\nimport org.springframework.context.annotation.Bean;\nimport org.springframework.context.annotation.Configuration;\nimport org.springframework.kafka.annotation.EnableKafka;\nimport org.springframework.kafka.core.DefaultKafkaProducerFactory;\nimport org.springframework.kafka.core.KafkaTemplate;\nimport org.springframework.kafka.core.ProducerFactory;\n\nimport java.util.HashMap;\nimport java.util.Map;\n\n@Configuration\n@EnableKafka\npublic class KafkaProducerConfig {\n\n    @Value(&quot;${kafka.producer.servers}&quot;)\n    private String servers;\n    @Value(&quot;${kafka.producer.retries}&quot;)\n    private int retries;\n    @Value(&quot;${kafka.producer.batch.size}&quot;)\n    private int batchSize;\n    @Value(&quot;${kafka.producer.linger}&quot;)\n    private int linger;\n    @Value(&quot;${kafka.producer.buffer.memory}&quot;)\n    private int bufferMemory;\n\n    public Map&lt;String, Object&gt; producerConfigs() {\n        Map&lt;String, Object&gt; props = new HashMap&lt;&gt;();\n        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, servers);\n        props.put(ProducerConfig.RETRIES_CONFIG, retries);\n        props.put(ProducerConfig.BATCH_SIZE_CONFIG, batchSize);\n        props.put(ProducerConfig.LINGER_MS_CONFIG, linger);\n        props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, bufferMemory);\n        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);\n        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);\n        return props;\n    }\n\n    public ProducerFactory&lt;String, String&gt; producerFactory() {\n        return new DefaultKafkaProducerFactory&lt;&gt;(producerConfigs());\n    }\n\n    @Bean\n    public KafkaTemplate&lt;String, String&gt; kafkaTemplate() {\n        return new KafkaTemplate&lt;String, String&gt;(producerFactory());\n    }\n}<\/code><\/pre>\n<p>\u6d4b\u8bd5\u6211\u4eec\u7684producer\uff0c\u5199\u4e00\u4e2aController\u3002\u5411<code>topic=test<\/code>\uff0c<code>key=key<\/code>\u53d1\u9001\u6d88\u606fmessage<\/p>\n<p>\u8bbf\u95ee\uff1a<a target=\"_blank\" rel=\"noopener\" href=\"http:\/\/localhost:8080\/kafka\/send?message=Joe.Ye\">http:\/\/localhost:8080\/kafka\/send?message=Joe.Ye<\/a><\/p>\n<pre><code class=\"language-java\">@RestController\n@RequestMapping(&quot;\/kafka&quot;)\npublic class ProducerController {\n\n    protected final Logger logger = LoggerFactory.getLogger(this.getClass());\n\n    @Autowired\n    private KafkaTemplate kafkaTemplate;\n\n    @ResponseBody\n    @RequestMapping(value = &quot;\/send&quot;, method = RequestMethod.GET)\n    public String sendKafka(HttpServletRequest request, HttpServletResponse response) {\n        try {\n            String message = request.getParameter(&quot;message&quot;);\n            logger.info(&quot;kafka\u7684\u6d88\u606f={}&quot;, message);\n            kafkaTemplate.send(&quot;test&quot;, &quot;key&quot;, message);\n            logger.info(&quot;\u53d1\u9001kafka\u6210\u529f&quot;);\n            return new Response(ResultCode.SUCCESS, &quot;\u53d1\u9001kafka\u6210\u529f&quot;, null).toString();\n        } catch (Exception e) {\n            logger.error(&quot;\u53d1\u9001kafka\u5931\u8d25&quot;, e);\n            return new Response(ResultCode.EXCEPTION, &quot;\u53d1\u9001kafka\u5931\u8d25&quot;, null).toString();\n        }\n    }\n\n}<\/code><\/pre>\n<h2>kafka consumer<\/h2>\n<p>\uff081\uff09\u901a\u8fc7<code>@Configuration<\/code>\u3001<code>@EnableKafka<\/code>\uff0c\u58f0\u660eConfig\u5e76\u4e14\u6253\u5f00KafkaTemplate\u80fd\u529b<br \/>\n\uff082\uff09\u901a\u8fc7<code>@Value<\/code>\u6ce8\u5165<code>application.properties<\/code>\u914d\u7f6e\u6587\u4ef6\u4e2d\u7684Kafka\u914d\u7f6e<br \/>\n\uff083\uff09\u901a\u8fc7<code>@Bean<\/code>\u751f\u6210bean<\/p>\n<pre><code class=\"language-java\">import me.yezhou.springboot.kafka.listener.Listener;\nimport org.apache.kafka.clients.consumer.ConsumerConfig;\nimport org.apache.kafka.common.serialization.StringDeserializer;\nimport org.springframework.beans.factory.annotation.Value;\nimport org.springframework.context.annotation.Bean;\nimport org.springframework.context.annotation.Configuration;\nimport org.springframework.kafka.annotation.EnableKafka;\nimport org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;\nimport org.springframework.kafka.config.KafkaListenerContainerFactory;\nimport org.springframework.kafka.core.ConsumerFactory;\nimport org.springframework.kafka.core.DefaultKafkaConsumerFactory;\nimport org.springframework.kafka.listener.ConcurrentMessageListenerContainer;\n\nimport java.util.HashMap;\nimport java.util.Map;\n\n@Configuration\n@EnableKafka\npublic class KafkaConsumerConfig {\n\n    @Value(&quot;${kafka.consumer.servers}&quot;)\n    private String servers;\n    @Value(&quot;${kafka.consumer.enable.auto.commit}&quot;)\n    private boolean enableAutoCommit;\n    @Value(&quot;${kafka.consumer.session.timeout}&quot;)\n    private String sessionTimeout;\n    @Value(&quot;${kafka.consumer.auto.commit.interval}&quot;)\n    private String autoCommitInterval;\n    @Value(&quot;${kafka.consumer.group.id}&quot;)\n    private String groupId;\n    @Value(&quot;${kafka.consumer.auto.offset.reset}&quot;)\n    private String autoOffsetReset;\n    @Value(&quot;${kafka.consumer.concurrency}&quot;)\n    private int concurrency;\n\n    @Bean\n    public KafkaListenerContainerFactory&lt;ConcurrentMessageListenerContainer&lt;String, String&gt;&gt; kafkaListenerContainerFactory() {\n        ConcurrentKafkaListenerContainerFactory&lt;String, String&gt; factory = new ConcurrentKafkaListenerContainerFactory&lt;&gt;();\n        factory.setConsumerFactory(consumerFactory());\n        factory.setConcurrency(concurrency);\n        factory.getContainerProperties().setPollTimeout(1500);\n        return factory;\n    }\n\n    public ConsumerFactory&lt;String, String&gt; consumerFactory() {\n        return new DefaultKafkaConsumerFactory&lt;&gt;(consumerConfigs());\n    }\n\n    public Map&lt;String, Object&gt; consumerConfigs() {\n        Map&lt;String, Object&gt; propsMap = new HashMap&lt;&gt;();\n        propsMap.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, servers);\n        propsMap.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, enableAutoCommit);\n        propsMap.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, autoCommitInterval);\n        propsMap.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, sessionTimeout);\n        propsMap.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);\n        propsMap.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);\n        propsMap.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);\n        propsMap.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, autoOffsetReset);\n        return propsMap;\n    }\n\n    @Bean\n    public Listener listener() {\n        return new Listener();\n    }\n}<\/code><\/pre>\n<p><code>new Listener()<\/code>\u751f\u6210\u4e00\u4e2abean\u7528\u6765\u5904\u7406\u4ecekafka\u8bfb\u53d6\u7684\u6570\u636e\u3002Listener\u7b80\u5355\u7684\u5b9e\u73b0\u5982\u4e0b\uff1a\u53ea\u662f\u7b80\u5355\u7684\u8bfb\u53d6\u5e76\u6253\u5370key\u548cmessage\u503c<\/p>\n<p><code>@KafkaListener<\/code>\u4e2dtopics\u5c5e\u6027\u7528\u4e8e\u6307\u5b9akafka topic\u540d\u79f0\uff0ctopic\u540d\u79f0\u7531\u6d88\u606f\u751f\u4ea7\u8005\u6307\u5b9a\uff0c\u4e5f\u5c31\u662f\u7531kafkaTemplate\u5728\u53d1\u9001\u6d88\u606f\u65f6\u6307\u5b9a\u3002<\/p>\n<pre><code class=\"language-java\">import org.apache.kafka.clients.consumer.ConsumerRecord;\nimport org.slf4j.Logger;\nimport org.slf4j.LoggerFactory;\nimport org.springframework.kafka.annotation.KafkaListener;\n\npublic class Listener {\n    protected final Logger logger = LoggerFactory.getLogger(this.getClass());\n\n    @KafkaListener(topics = {&quot;test&quot;})\n    public void listen(ConsumerRecord&lt;?, ?&gt; record) {\n        logger.info(&quot;Kafka\u7684key: &quot; + record.key());\n        logger.info(&quot;Kafka\u7684value: &quot; + record.value().toString());\n    }\n}<\/code><\/pre>\n<pre><code>Kafka version : 2.0.1\nKafka commitId : fa14705e51bd2ce5\nCluster ID: HhhOI8SET2yHTesb9OQKbg\n\u53d1\u9001kafka\u6210\u529f\nKafka\u7684key: key\nKafka\u7684value: Joe.Ye<\/code><\/pre>\n<h2>\u6ce8\u610f\u4e8b\u9879<\/h2>\n<p>\uff081\uff09\u914d\u7f6eKafka\u65f6\u6700\u597d\u7528\u5b8c\u5168bind\u7f51\u7edcip\u7684\u65b9\u5f0f\uff0c\u800c\u4e0d\u662flocalhost\u6216\u8005127.0.0.1<br \/>\n\uff082\uff09\u6700\u597d\u4e0d\u8981\u4f7f\u7528Kafka\u81ea\u5e26\u7684ZooKeeper\u90e8\u7f72Kafka\uff0c\u53ef\u80fd\u5bfc\u81f4\u8bbf\u95ee\u4e0d\u901a<br \/>\n\uff083\uff09\u5b9a\u4e49\u76d1\u542c\u6d88\u606f\u914d\u7f6e\u65f6\uff0cGROUP_ID_CONFIG\u914d\u7f6e\u9879\u7684\u503c\u7528\u4e8e\u6307\u5b9a\u6d88\u8d39\u8005\u7ec4\u7684\u540d\u79f0\uff0c\u5982\u679c\u540c\u7ec4\u4e2d\u5b58\u5728\u591a\u4e2a\u76d1\u542c\u5668\u5bf9\u8c61\u5219\u53ea\u6709\u4e00\u4e2a\u76d1\u542c\u5668\u5bf9\u8c61\u80fd\u6536\u5230\u6d88\u606f<\/p>\n","protected":false},"excerpt":{"rendered":"<p>\u4ecb\u7ecd\u5982\u4f55\u5728Spring Boot\u9879\u76ee\u4e2d\u96c6\u6210Kafka\u6536\u53d1Message pom.xml\u4f9d\u8d56 &lt;!&#8211; h [&hellip;]<\/p>\n","protected":false},"author":1,"featured_media":0,"comment_status":"open","ping_status":"open","sticky":false,"template":"","format":"standard","meta":{"footnotes":""},"categories":[31,41],"tags":[183],"class_list":["post-872","post","type-post","status-publish","format-standard","hentry","category-kafka","category-spring-boot","tag-kafka"],"_links":{"self":[{"href":"https:\/\/www.appblog.cn\/index.php\/wp-json\/wp\/v2\/posts\/872","targetHints":{"allow":["GET"]}}],"collection":[{"href":"https:\/\/www.appblog.cn\/index.php\/wp-json\/wp\/v2\/posts"}],"about":[{"href":"https:\/\/www.appblog.cn\/index.php\/wp-json\/wp\/v2\/types\/post"}],"author":[{"embeddable":true,"href":"https:\/\/www.appblog.cn\/index.php\/wp-json\/wp\/v2\/users\/1"}],"replies":[{"embeddable":true,"href":"https:\/\/www.appblog.cn\/index.php\/wp-json\/wp\/v2\/comments?post=872"}],"version-history":[{"count":0,"href":"https:\/\/www.appblog.cn\/index.php\/wp-json\/wp\/v2\/posts\/872\/revisions"}],"wp:attachment":[{"href":"https:\/\/www.appblog.cn\/index.php\/wp-json\/wp\/v2\/media?parent=872"}],"wp:term":[{"taxonomy":"category","embeddable":true,"href":"https:\/\/www.appblog.cn\/index.php\/wp-json\/wp\/v2\/categories?post=872"},{"taxonomy":"post_tag","embeddable":true,"href":"https:\/\/www.appblog.cn\/index.php\/wp-json\/wp\/v2\/tags?post=872"}],"curies":[{"name":"wp","href":"https:\/\/api.w.org\/{rel}","templated":true}]}}