{"id":1124,"date":"2020-05-31T15:01:29","date_gmt":"2020-05-31T07:01:29","guid":{"rendered":"http:\/\/www.rain1024.com\/?p=1124"},"modified":"2023-08-07T20:57:24","modified_gmt":"2023-08-07T12:57:24","slug":"%e4%bd%bf%e7%94%a8%e4%b8%a4%e7%a7%8d%e5%a4%9a%e7%ba%bf%e7%a8%8b%e6%a8%a1%e5%bc%8f%e6%b6%88%e8%b4%b9%e6%95%b0%e6%8d%ae","status":"publish","type":"post","link":"http:\/\/rain1024.com\/index.php\/2020\/05\/31\/%e4%bd%bf%e7%94%a8%e4%b8%a4%e7%a7%8d%e5%a4%9a%e7%ba%bf%e7%a8%8b%e6%a8%a1%e5%bc%8f%e6%b6%88%e8%b4%b9%e6%95%b0%e6%8d%ae\/","title":{"rendered":"\u4f7f\u7528\u4e24\u79cd\u591a\u7ebf\u7a0b\u6a21\u5f0f\u6d88\u8d39\u6570\u636e"},"content":{"rendered":"<h1>\u4f7f\u7528\u4e24\u79cd\u591a\u7ebf\u7a0b\u6a21\u5f0f\u6d88\u8d39\u6570\u636e<\/h1>\n<p>KafkaProducer\u662f\u7ebf\u7a0b\u5b89\u5168\u7684,\u7136\u800c KafkaConsumer\u5374\u662f\u975e\u7ebf\u7a0b\u5b89\u5168\u7684\u3002 Kafka Consumer\u4e2d\u5b9a\u4e49\u4e86\u4e00\u4e2a acquire(\u65b9\u6cd5,\u7528\u6765\u68c0\u6d4b\u5f53\u524d\u662f\u5426\u53ea\u6709\u4e00\u4e2a\u7ebf\u7a0b\u5728\u64cd\u4f5c,\u82e5\u6709\u5176\u4ed6\u7ebf\u7a0b\u6b63\u5728\u64cd\u4f5c\u5219\u4f1a\u629b\u51fa Concurrentmodifcationexception\u5f02\u5e38:<\/p>\n<p>java.util.ConcurrentModificationException: KafkaConsumer is not safe for multi-threaded access.<\/p>\n<p>KafkaConsumer\u975e\u7ebf\u7a0b\u5b89\u5168\u5e76\u4e0d\u610f\u5473\u7740\u6211\u4eec\u5728\u6d88\u8d39\u6d88\u606f\u7684\u65f6\u5019\u53ea\u80fd\u4ee5\u5355\u7ebf\u7a0b\u7684\u65b9\u5f0f\u6267\u884c\u3002\u5982\u679c\u751f\u4ea7\u8005\u53d1\u9001\u6d88\u606f\u7684\u901f\u5ea6\u5927\u4e8e\u6d88\u8d39\u8005\u5904\u7406\u6d88\u606f\u7684\u901f\u5ea6,\u90a3\u4e48\u5c31\u4f1a\u6709\u8d8a\u6765\u8d8a\u591a\u7684\u6d88\u606f\u5f97\u4e0d\u5230\u53ca\u65f6\u7684\u6d88\u8d39,\u9020\u6210\u4e86\u4e00\u5b9a\u7684\u5ef6\u8fdf\u3002\u9664\u6b64\u4e4b\u5916,\u7531\u4e8eKafka\u4e2d\u6d88\u606f\u4fdd\u7559\u673a\u5236\u7684\u4f5c\u7528,\u6709\u4e9b\u6d88\u606f\u6709\u53ef\u80fd\u5728\u88ab\u6d88\u8d39\u4e4b\u524d\u5c31\u88ab\u6e05\u7406\u4e86,\u4ece\u800c\u9020\u6210\u6d88\u606f\u7684\u4e22\u5931\u3002\u6211\u4eec\u53ef\u4ee5\u901a\u8fc7\u591a\u7ebf\u7a0b\u7684\u65b9\u5f0f\u6765\u5b9e\u73b0\u6d88\u606f\u6d88\u8d39,\u591a\u7ebf\u7a0b\u7684\u76ee\u7684\u5c31\u662f\u4e3a\u4e86\u63d0\u9ad8\u6574\u4f53\u7684\u6d88\u8d39\u80fd\u529b\u3002\u591a\u7ebf\u7a0b\u7684\u5b9e\u73b0\u65b9\u5f0f\u6709\u591a\u79cd,\u7b2c\u4e00\u79cd\u4e5f\u662f\u6700\u5e38\u89c1\u7684\u65b9\u5f0f:\u7ebf\u7a0b\u5c01\u95ed,\u5373\u4e3a\u6bcf\u4e2a\u7ebf\u7a0b\u5b9e\u4f8b\u5316\u4e00\u4e2aKafkaConsumer\u5bf9\u8c61,\u5982\u56fe3-10\u6240\u793a\u3002<\/p>\n<p><img decoding=\"async\" src=\"http:\/\/uos.rain1024.com\/image\/image-20200531143640944.png\" alt=\"image-20200531143640944\" \/><\/p>\n<h3>\u7b2c\u4e00\u79cd\u591a\u7ebf\u7a0b\u6d88\u8d39\u5b9e\u73b0\u65b9\u5f0f<\/h3>\n<p>\u4e00\u4e2a\u7ebf\u7a0b\u5bf9\u5e94\u4e00\u4e2aKafkaConsumer\u5b9e\u4f8b,\u6211\u4eec\u53ef\u4ee5\u79f0\u4e4b\u4e3a\u6d88\u8d39\u7ebf\u7a0b\u3002\u4e00\u4e2a\u6d88\u8d39\u7ebf\u7a0b\u53ef\u4ee5\u6d88\u8d39\u4e00\u4e2a\u6216\u591a\u4e2a\u5206\u533a\u4e2d\u7684\u6d88\u606f,\u6240\u6709\u7684\u6d88\u8d39\u7ebf\u7a0b\u90fd\u96b6\u5c5e\u4e8e\u540c\u4e00\u4e2a\u6d88\u8d39\u7ec4\u3002\u8fd9\u79cd\u5b9e\u73b0\u65b9\u5f0f\u7684\u5e76\u53d1\u5ea6\u53d7\u9650\u4e8e\u5206\u533a\u7684\u5b9e\u9645\u4e2a\u6570,\u5f53\u6d88\u8d39\u7ebf\u7a0b\u7684\u4e2a\u6570\u5927\u4e8e\u5206\u533a\u6570\u65f6,\u5c31\u6709\u90e8\u5206\u6d88\u8d39\u7ebf\u7a0b\u4e00\u76f4\u5904\u4e8e\u7a7a\u95f2\u7684\u72b6\u6001\u3002<\/p>\n<pre><code>package com.rain.demo;\n\nimport org.apache.kafka.clients.consumer.ConsumerConfig;\nimport org.apache.kafka.clients.consumer.ConsumerRecord;\nimport org.apache.kafka.clients.consumer.ConsumerRecords;\nimport org.apache.kafka.clients.consumer.KafkaConsumer;\nimport org.apache.kafka.common.serialization.StringDeserializer;\n\nimport java.lang.reflect.Array;\nimport java.time.Duration;\nimport java.util.Arrays;\nimport java.util.Properties;\n\n\/**\n * @Author: wcy\n * @Date: 2020\/5\/31\n *\/\npublic class FirstMultiConsumerThreadDemo {\n\n    public static final String brokerList = \"nas-cluster1:9092\";\n    public static final String topic = \"test.topic\";\n    public static final String groupId = \"group.demo\";\n\n    public static Properties initConfig(){\n        Properties properties = new Properties();\n        properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, brokerList);\n        properties.put(ConsumerConfig.GROUP_ID_CONFIG,groupId);\n        properties.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG,true);\n        properties.put(\"key.deserializer\", StringDeserializer.class.getName());\n        properties.put(\"value.deserializer\",StringDeserializer.class.getName());\n        return properties;\n    }\n\n    public static void main(String[] args) {\n        Properties props = initConfig();\n        int consumerThreadNum = 4;\n        for (int i = 0; i &lt; consumerThreadNum; i++) {\n            new KafkaConsumerThread(props,topic).start();\n        }\n    }\n\n    public static class KafkaConsumerThread extends Thread{\n        private KafkaConsumer&lt;String,String&gt; kafkaConsumer;\n\n\n        public KafkaConsumerThread(Properties props, String topic) {\n            this.kafkaConsumer = new KafkaConsumer&lt;&gt;(props);\n            this.kafkaConsumer.subscribe(Arrays.asList(topic));\n        }\n\n        @Override\n        public void run() {\n            try {\n                while (true){\n                    ConsumerRecords&lt;String,String&gt; records = kafkaConsumer.poll(Duration.ofMillis(100));\n                    for (ConsumerRecord&lt;String,String&gt; record : records){\n                        \/\/\u5b9e\u73b0\u5904\u7406\u903b\u8f91\n                        System.out.println(record.value());\n                    }\n                }\n            }catch (Exception e){\n                e.printStackTrace();\n            }finally {\n                kafkaConsumer.close();\n            }\n        }\n    }\n}\n<\/code><\/pre>\n<p>\u5185\u90e8\u7c7b Kafka Consumer Thread\u4ee3\u8868\u6d88\u8d39\u7ebf\u7a0b,\u5176\u5185\u90e8\u5305\u88f9\u7740\u4e00\u4e2a\u72ec\u7acb\u7684 Kafka Consumer\u5b9e\u4f8b\u3002\u901a\u8fc7\u5916\u90e8\u7c7b\u7684 maino\u65b9\u6cd5\u6765\u542f\u52a8\u591a\u4e2a\u6d88\u8d39\u7ebf\u7a0b,\u6d88\u8d39\u7ebf\u7a0b\u7684\u6570\u91cf\u7531 consumer Threadnum\u53d8\u91cf\u6307\u5b9a\u3002\u4e00\u822c\u4e00\u4e2a\u4e3b\u9898\u7684\u5206\u533a\u6570\u4e8b\u5148\u53ef\u4ee5\u77e5\u6653,\u53ef\u4ee5\u5c06 consumer Threadnum\u8bbe\u7f6e\u6210\u4e0d\u5927\u4e8e\u5206\u533a\u6570\u7684\u503c,\u5982\u679c\u4e0d\u77e5\u9053\u4e3b\u9898\u7684\u5206\u533a\u6570,\u90a3\u4e48\u4e5f\u53ef\u4ee5\u901a\u8fc7 Kafka Consumer\u7c7b\u7684 partitionsforo\u65b9\u6cd5\u6765\u95f4\u63a5\u83b7\u53d6,\u8fdb\u800c\u518d\u8bbe\u7f6e\u5408\u7406\u7684 consumer Threadnum\u503c\u3002<br \/>\n\u4e0a\u9762\u8fd9\u79cd\u591a\u7ebf\u7a0b\u7684\u5b9e\u73b0\u65b9\u5f0f\u548c\u5f00\u542f\u591a\u4e2a\u6d88\u8d39\u8fdb\u7a0b\u7684\u65b9\u5f0f\u6ca1\u6709\u672c\u8d28\u4e0a\u7684\u533a\u522b,\u5b83\u7684\u4f18\u70b9\u662f\u6bcf\u4e2a\u7ebf\u7a0b\u53ef\u4ee5\u6309\u987a\u5e8f\u6d88\u8d39\u5404\u4e2a\u5206\u533a\u4e2d\u7684\u6d88\u606f\u3002\u7f3a\u70b9\u4e5f\u5f88\u660e\u663e,\u6bcf\u4e2a\u6d88\u8d39\u7ebf\u7a0b\u90fd\u8981\u7ef4\u62a4\u4e00\u4e2a\u72ec\u7acb\u7684TCP\u8fde\u63a5,\u5982\u679c\u5206\u533a\u6570\u548c consumer Threadnum\u7684\u503c\u90fd\u5f88\u5927,\u90a3\u4e48\u4f1a\u9020\u6210\u4e0d\u5c0f\u7684\u7cfb\u7edf\u5f00\u9500\u3002<\/p>\n<h3>\u7b2c\u4e8c\u79cd\u57fa\u4e8e\u6570\u636e\u5904\u7406\u7684\u591a\u7ebf\u7a0b\u6d88\u8d39\u5b9e\u73b0<\/h3>\n<p>\u5982\u679c\u5904\u7406\u6570\u636e\u7684\u5730\u65b9\u5bf9\u6d88\u606f\u7684\u5904\u7406\u975e\u5e38\u8fc5\u901f,\u90a3\u4e48pollo\u62c9\u53d6\u7684\u9891\u6b21\u4e5f\u4f1a\u66f4\u9ad8,\u8fdb\u800c\u6574\u4f53\u6d88\u8d39\u7684\u6027\u80fd\u4e5f\u4f1a\u63d0\u5347;\u76f8\u53cd,\u5982\u679c\u5728\u8fd9\u91cc\u5bf9\u6d88\u606f\u7684\u5904\u7406\u7f13\u6162,\u6bd4\u5982\u8fdb\u884c\u4e00\u4e2a\u4e8b\u52a1\u6027\u64cd\u4f5c,\u6216\u8005\u7b49\u5f85\u4e00\u4e2aRPC\u7684\u540c\u6b65\u54cd\u5e94,\u90a3\u4e48poll(\u62c9\u53d6\u7684\u9891\u6b21\u4e5f\u4f1a\u968f\u4e4b\u4e0b\u964d,\u8fdb\u800c\u9020\u6210\u6574\u4f53\u6d88\u8d39\u6027\u80fd\u7684\u4e0b\u964d\u3002\u4e00\u822c\u800c\u8a00, pol()\u62c9\u53d6\u6d88\u606f\u7684\u901f\u5ea6\u662f\u76f8\u5f53\u5feb\u7684,\u800c\u6574\u4f53\u6d88\u8d39\u7684\u74f6\u9888\u4e5f\u6b63\u662f\u5728\u5904\u7406\u6d88\u606f\u8fd9\u4e00\u5757,\u5982\u679c\u6211\u4eec\u901a\u8fc7\u4e00\u5b9a\u7684\u65b9\u5f0f\u6765\u6539\u8fdb\u8fd9\u4e00\u90e8\u5206,\u90a3\u4e48\u6211\u4eec\u5c31\u80fd\u5e26\u52a8\u6574\u4f53\u6d88\u8d39\u6027\u80fd\u7684\u63d0\u5347\uff0c\u56e0\u6b64\u5c06\u5904\u7406\u6d88\u606f\u6a21\u5757\u6539\u6210\u591a\u7ebf\u7a0b\u7684\u5b9e\u73b0\u65b9\u5f0f\u3002<\/p>\n<p><img decoding=\"async\" src=\"http:\/\/uos.rain1024.com\/image\/image-20200531145522193.png\" alt=\"image-20200531145522193\" \/><\/p>\n<pre><code>package com.rain.demo;\n\nimport org.apache.kafka.clients.consumer.ConsumerConfig;\nimport org.apache.kafka.clients.consumer.ConsumerRecord;\nimport org.apache.kafka.clients.consumer.ConsumerRecords;\nimport org.apache.kafka.clients.consumer.KafkaConsumer;\nimport org.apache.kafka.common.serialization.StringDeserializer;\n\nimport java.time.Duration;\nimport java.util.Collections;\nimport java.util.Properties;\nimport java.util.concurrent.ArrayBlockingQueue;\nimport java.util.concurrent.ExecutorService;\nimport java.util.concurrent.ThreadPoolExecutor;\nimport java.util.concurrent.TimeUnit;\n\n\/**\n * @Author: wcy\n * @Date: 2020\/5\/31\n *\/\npublic class SecondMultiConsumerThreadDemo {\n    public static final String brokerList = \"nas-cluster1:9092\";\n    public static final String topic = \"test.topic\";\n    public static final String groupId = \"group.demo\";\n\n    public static Properties initConfig(){\n        Properties properties = new Properties();\n        properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, brokerList);\n        properties.put(ConsumerConfig.GROUP_ID_CONFIG,groupId);\n        properties.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG,true);\n        properties.put(\"key.deserializer\", StringDeserializer.class.getName());\n        properties.put(\"value.deserializer\",StringDeserializer.class.getName());\n        return properties;\n    }\n\n    public static void main(String[] args) {\n        Properties properties = initConfig();\n        KafkaConsumerThread consumerThread = new KafkaConsumerThread(properties,topic,\n                Runtime.getRuntime().availableProcessors());\n        consumerThread.start();\n    }\n\n    public static class KafkaConsumerThread extends Thread{\n        private KafkaConsumer&lt;String,String&gt; kafkaConsumer;\n        private ExecutorService executorService;\n        private int threadNumber;\n\n\n        public KafkaConsumerThread(Properties properties, String topic, int availableProcessors) {\n            kafkaConsumer = new KafkaConsumer&lt;String, String&gt;(properties);\n            kafkaConsumer.subscribe(Collections.singletonList(topic));\n            this.threadNumber = availableProcessors;\n            executorService = new ThreadPoolExecutor(threadNumber,threadNumber,0L, TimeUnit.MILLISECONDS,\n                    new ArrayBlockingQueue&lt;&gt;(1000),new ThreadPoolExecutor.CallerRunsPolicy());\n\n        }\n\n        @Override\n        public void run() {\n            try {\n                while (true){\n                    ConsumerRecords&lt;String,String&gt; records = kafkaConsumer.poll(Duration.ofMillis(100));\n                    if (!records.isEmpty()){\n                        executorService.submit(new RecordsHandler(records));\n                    }\n                }\n            }catch (Exception e){\n                e.printStackTrace();\n            }finally {\n                kafkaConsumer.close();\n            }\n        }\n    }\n    public static class RecordsHandler extends Thread{\n        public final ConsumerRecords&lt;String,String&gt; records;\n\n        public RecordsHandler(ConsumerRecords&lt;String, String&gt; records) {\n            this.records = records;\n        }\n\n        @Override\n        public void run() {\n            for (ConsumerRecord&lt;String,String&gt; record : records){\n                \/\/\u5b9e\u73b0\u5904\u7406\u903b\u8f91\n                System.out.println(record.value());\n            }\n        }\n    }\n}\n<\/code><\/pre>\n<p>\u4ee3\u7801\u4e2d Recordhandler\u7c7b\u662f\u7528\u6765\u5904\u7406\u6d88\u606f\u7684,\u800c Kafka Thread\u7c7b\u5bf9\u5e94\u7684\u662f\u4e00\u4e2a\u6d88\u8d39\u7ebf\u7a0b,\u91cc\u9762\u901a\u8fc7\u7ebf\u7a0b\u6c60\u7684\u65b9\u5f0f\u6765\u8c03\u7528 Recordhandler\u5904\u7406\u4e00\u6279\u6279\u7684\u6d88\u606f\u3002\u6ce8\u610fKafka Consumer Thread\u7c7b\u4e2d Threadpoolexecutor\u91cc\u7684\u6700\u540e\u4e00\u4e2a\u53c2\u6570\u8bbe\u7f6e\u7684\u662f Callerrunspolicyo, \u8fd9\u6837\u53ef\u4ee5\u9632\u6b62\u7ebf\u7a0b\u6c60\u7684\u603b\u4f53\u6d88\u8d39\u80fd\u529b\u8ddf\u4e0d\u4e0apolO\u62c9\u53d6\u7684\u80fd\u529b,\u4ece\u800c\u5bfc\u81f4\u5f02\u5e38\u73b0\u8c61\u7684\u53d1\u751f\u3002\u7b2c\u4e09\u79cd\u5b9e\u73b0\u65b9\u5f0f\u8fd8\u53ef\u4ee5\u6a2a\u5411\u6269\u5c55,\u901a\u8fc7\u5f00\u542f\u591a\u4e2a Kafka Consumerthread\u5b9e\u4f8b\u6765\u8fdb\u4e00\u6b65\u63d0\u5347\u6574\u4f53\u7684\u6d88\u8d39\u80fd\u529b\u3002<\/p>\n","protected":false},"excerpt":{"rendered":"<p>\u4f7f\u7528\u4e24\u79cd\u591a\u7ebf\u7a0b\u6a21\u5f0f\u6d88\u8d39\u6570\u636e KafkaProducer\u662f\u7ebf\u7a0b\u5b89\u5168\u7684,\u7136\u800c KafkaConsumer\u5374\u662f\u975e\u7ebf\u7a0b\u2026 <span class=\"read-more\"><a href=\"http:\/\/rain1024.com\/index.php\/2020\/05\/31\/%e4%bd%bf%e7%94%a8%e4%b8%a4%e7%a7%8d%e5%a4%9a%e7%ba%bf%e7%a8%8b%e6%a8%a1%e5%bc%8f%e6%b6%88%e8%b4%b9%e6%95%b0%e6%8d%ae\/\">Read More &raquo;<\/a><\/span><\/p>\n","protected":false},"author":1,"featured_media":0,"comment_status":"open","ping_status":"open","sticky":false,"template":"","format":"standard","meta":{"footnotes":""},"categories":[7],"tags":[],"class_list":["post-1124","post","type-post","status-publish","format-standard","hentry","category-kafka-hadoop"],"_links":{"self":[{"href":"http:\/\/rain1024.com\/index.php\/wp-json\/wp\/v2\/posts\/1124","targetHints":{"allow":["GET"]}}],"collection":[{"href":"http:\/\/rain1024.com\/index.php\/wp-json\/wp\/v2\/posts"}],"about":[{"href":"http:\/\/rain1024.com\/index.php\/wp-json\/wp\/v2\/types\/post"}],"author":[{"embeddable":true,"href":"http:\/\/rain1024.com\/index.php\/wp-json\/wp\/v2\/users\/1"}],"replies":[{"embeddable":true,"href":"http:\/\/rain1024.com\/index.php\/wp-json\/wp\/v2\/comments?post=1124"}],"version-history":[{"count":1,"href":"http:\/\/rain1024.com\/index.php\/wp-json\/wp\/v2\/posts\/1124\/revisions"}],"predecessor-version":[{"id":1357,"href":"http:\/\/rain1024.com\/index.php\/wp-json\/wp\/v2\/posts\/1124\/revisions\/1357"}],"wp:attachment":[{"href":"http:\/\/rain1024.com\/index.php\/wp-json\/wp\/v2\/media?parent=1124"}],"wp:term":[{"taxonomy":"category","embeddable":true,"href":"http:\/\/rain1024.com\/index.php\/wp-json\/wp\/v2\/categories?post=1124"},{"taxonomy":"post_tag","embeddable":true,"href":"http:\/\/rain1024.com\/index.php\/wp-json\/wp\/v2\/tags?post=1124"}],"curies":[{"name":"wp","href":"https:\/\/api.w.org\/{rel}","templated":true}]}}