Kafka生产者发送消息的三种方式

Kafka是一种分布式的基于发布/订阅的消息系统,它的高吞吐量、灵活的offset是其它消息系统所没有的。

Kafka发送消息主要有三种方式:

1.发送并忘记 2.同步发送 3.异步发送+回调函数

下面以单节点的方式分别用三种方法发送1w条消息测试:

方式一:发送并忘记(不关心消息是否正常到达,对返回结果不做任何判断处理)

发送并忘记的方式本质上也是一种异步的方式,只是它不会获取消息发送的返回结果,这种方式的吞吐量是最高的,但是无法保证消息的可靠性:

1 1 import pickle 2 2 import time 3 3 from kafka import KafkaProducer 4 4 5 5 producer = KafkaProducer(bootstrap_servers=['192.168.33.11:9092'], 6 6 key_serializer=lambda k: pickle.dumps(k), 7 7 value_serializer=lambda v: pickle.dumps(v)) 8 8 9 9 start_time = time.time() 1010 for i in range(0, 10000): 1111 print('------{}---------'.format(i)) 1212 future = producer.send('test_topic', key='num', value=i, partition=0) 1313 1414 # 将缓冲区的全部消息push到broker当中 1515 producer.flush() 1616 producer.close() 1717 1818 end_time = time.time() 1919 time_counts = end_time - start_time 2020 print(time_counts)

测试结果:1.88s

方式二:同步发送(通过get方法等待Kafka的响应,判断消息是否发送成功)

以同步的方式发送消息时,一条一条的发送,对每条消息返回的结果判断, 可以明确地知道每条消息的发送情况,但是由于同步的方式会阻塞,只有当消息通过get返回future对象时,才会继续下一条消息的发送:

1 1 import pickle 2 2 import time 3 3 from kafka import KafkaProducer 4 4 from kafka.errors import kafka_errors 5 5 6 6 producer = KafkaProducer( 7 7 bootstrap_servers=['192.168.33.11:9092'], 8 8 key_serializer=lambda k: pickle.dumps(k), 9 9 value_serializer=lambda v: pickle.dumps(v) 1010 ) 1111 1212 start_time = time.time() 1313 for i in range(0, 10000): 1414 print('------{}---------'.format(i)) 1515 future = producer.send(topic="test_topic", key="num", value=i) 1616 # 同步阻塞,通过调用get()方法进而保证一定程序是有序的. 1717 try: 1818 record_metadata = future.get(timeout=10) 1919 # print(record_metadata.topic) 2020 # print(record_metadata.partition) 2121 # print(record_metadata.offset) 2222 except kafka_errors as e: 2323 print(str(e)) 2424 2525 end_time = time.time() 2626 time_counts = end_time - start_time 2727 print(time_counts)

测试结果:16s

方式三:异步发送+回调函数(消息以异步的方式发送,通过回调函数返回消息发送成功/失败)

在调用send方法发送消息的同时,指定一个回调函数,服务器在返回响应时会调用该回调函数,通过回调函数能够对异常情况进行处理,当调用了回调函数时,只有回调函数执行完毕生产者才会结束,否则一直会阻塞:

1 1 import pickle 2 2 import time 3 3 from kafka import KafkaProducer 4 4 5 5 producer = KafkaProducer( 6 6 bootstrap_servers=['192.168.33.11:9092'], 7 7 key_serializer=lambda k: pickle.dumps(k), 8 8 value_serializer=lambda v: pickle.dumps(v) 9 9 ) 1010 1111 1212 def on_send_success(*args, **kwargs): 1313 """ 1414 发送成功的回调函数 1515 :param args: 1616 :param kwargs: 1717 :return: 1818 """ 1919 return args 2020 2121 2222 def on_send_error(*args, **kwargs): 2323 """ 2424 发送失败的回调函数 2525 :param args: 2626 :param kwargs: 2727 :return: 2828 """ 2929 3030 return args 3131 3232 3333 start_time = time.time() 3434 for i in range(0, 10000): 3535 print('------{}---------'.format(i)) 3636 # 如果成功,传进record_metadata,如果失败,传进Exception. 3737 producer.send( 3838 topic="test_topic", key="num", value=i 3939 ).add_callback(on_send_success).add_errback(on_send_error) 4040 4141 producer.flush() 4242 producer.close() 4343 4444 end_time = time.time() 4545 time_counts = end_time - start_time 4646 print(time_counts)

测试结果:2.15s

三种方式虽然在时间上有所差别,但并不是说时间越快的越好,具体要看业务的应用场景:

场景1:如果业务要求消息必须是按顺序发送的,那么可以使用同步的方式,并且只能在一个partation上,结合参数设置retries的值让发送失败时重试,设置max_in_flight_requests_per_connection=1,可以控制生产者在收到服务器晌应之前只能发送1个消息,从而控制消息顺序发送;

场景2:如果业务只关心消息的吞吐量,容许少量消息发送失败,也不关注消息的发送顺序,那么可以使用发送并忘记的方式,并配合参数acks=0,这样生产者不需要等待服务器的响应,以网络能支持的最大速度发送消息;

场景3:如果业务需要知道消息发送是否成功,并且对消息的顺序不关心,那么可以用异步+回调的方式来发送消息,配合参数retries=0,并将发送失败的消息记录到日志文件中;

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

thinkcmf+jsapi 实现微信支付

首先从小程序端接收订单号、金额等参数,然后后台进行统一下单,把微信支付的订单号返回,在把订单号发送给前台,前台拉起支付,返回参数后更改支付状态。。。回调publicfunctionnotify(){$wechatDb::name('wechat')where('status',1)find();

Kafka概述及安装部署

一、Kafka概述1.Kafka是一个分布式流媒体平台,它有三个关键功能:(1)发布和订阅记录流,类似于消息队列或企业消息传递系统;(2)以容错的持久方式存储记录流;(3)记录发送时处理流。2.Kafka通常应用的两大类应用(1)构建在系统或应用程序之间的可靠获取数据的实时流数据管道;(2)构建转换或响应数据流的实施

FusionInsight大数据开发

Kafka应用开发1.了解Kafka应用开发适用场景2.熟悉Kafka应用开发流程3.熟悉并使用Kafka常用API4.进行Kafka应用开发Kafka的定义Kafka是一个高吞吐、分布式、基于发布订阅的消息系统Kafka有如下几个特点:1.高吞吐量2.消息持久化到磁