Pub/Sub回调写入MySQL出现连接错误与内存释放错误如何解决
问题根因
- MySQL连接对象线程不安全:Google Cloud Pub/Sub 订阅的回调函数默认在多线程环境下并发执行,全局共用同一个连接对象时,多个线程同时操作连接的底层内存结构,会触发C扩展层的内存访问错误,也就是你遇到的malloc崩溃报错;同时并发操作也会破坏连接状态,导致连接不可用。
- 全局连接无保活/重连逻辑:MySQL服务端默认会主动断开闲置超过
wait_timeout(默认8小时)的连接,闲置后再使用就会抛出连接丢失的错误。 - 提前ack消息存在隐患:你在处理消息前就执行了ack,若数据库写入失败,消息不会重试,会直接丢失。
解决方案示例
新手友好版(低并发场景适用)
每次回调独立创建数据库连接,避免多线程共享冲突,逻辑简单不易出错,适合消息量不大的场景。
from google.cloud import pubsub_v1 import mysql.connector from mysql.connector import Error import json subscriber = pubsub_v1.SubscriberClient() publisher = pubsub_v1.PublisherClient() subscription_path = subscriber.subscription_path('sim', 'sub') topic_path = publisher.topic_path('sim', 'topic') DB_CONFIG = { "host": "localhost", "user": "root", "password": "pass", "database": "db", "auth_plugin": "mysql_native_password", } def callback(message): print(f'Received {message.data}') cnx = None cursor = None try: # 解析消息数据 data = json.loads(message.data) # 本次回调独立创建数据库连接 cnx = mysql.connector.connect(**DB_CONFIG) cursor = cnx.cursor() insert_record = """ INSERT INTO Device (deviceId, temperature, location, time) VALUES (%s, %s, point(%s,%s), %s) """ record = ( data['deviceId'], data['temperature'], data['Location']['lat'], data['Location']['long'], data['time'] ) cursor.execute(insert_record, record) cnx.commit() # 写入成功再确认消息 message.ack() print(f'Acknowledged {message.message_id}') except Exception as e: print(f"处理消息{message.message_id}失败: {str(e)}") # 处理失败,消息重新入队重试 message.nack() finally: # 释放资源 if cursor: cursor.close() if cnx and cnx.is_connected(): cnx.close() # 订阅不存在时才创建,重复创建会报错,建议提前在控制台创建好订阅去掉这行 try: subscriber.create_subscription(name = subscription_path, topic = topic_path) except Exception as e: print(f"订阅已存在或创建失败: {str(e)}") streaming_pull_future = subscriber.subscribe(subscription_path, callback=callback) print(f'Listening for messages on {subscription_path}') with subscriber: try: streaming_pull_future.result() except TimeoutError: streaming_pull_future.cancel() streaming_pull_future.result()
高并发优化版(连接池实现)
消息量较大时使用连接池复用连接,避免频繁创建销毁连接的开销:
from google.cloud import pubsub_v1 import mysql.connector from mysql.connector import Error, pooling import json subscriber = pubsub_v1.SubscriberClient() publisher = pubsub_v1.PublisherClient() subscription_path = subscriber.subscription_path('sim', 'sub') topic_path = publisher.topic_path('sim', 'topic') DB_CONFIG = { "host": "localhost", "user": "root", "password": "pass", "database": "db", "auth_plugin": "mysql_native_password", } # 初始化连接池,pool_size大小根据实际并发量调整,最大不超过Pub/Sub订阅的最大并发数 cnx_pool = pooling.MySQLConnectionPool(pool_name="device_pool", pool_size=5, **DB_CONFIG) def callback(message): print(f'Received {message.data}') cnx = None cursor = None try: data = json.loads(message.data) # 从连接池获取连接 cnx = cnx_pool.get_connection() cursor = cnx.cursor() insert_record = """ INSERT INTO Device (deviceId, temperature, location, time) VALUES (%s, %s, point(%s,%s), %s) """ record = ( data['deviceId'], data['temperature'], data['Location']['lat'], data['Location']['long'], data['time'] ) cursor.execute(insert_record, record) cnx.commit() message.ack() print(f'Acknowledged {message.message_id}') except Exception as e: print(f"处理消息{message.message_id}失败: {str(e)}") message.nack() finally: if cursor: cursor.close() if cnx and cnx.is_connected(): # 连接池模式下调用close会自动将连接放回池中 cnx.close() try: subscriber.create_subscription(name = subscription_path, topic = topic_path) except Exception as e: print(f"订阅已存在或创建失败: {str(e)}") streaming_pull_future = subscriber.subscribe(subscription_path, callback=callback) print(f'Listening for messages on {subscription_path}') with subscriber: try: streaming_pull_future.result() except TimeoutError: streaming_pull_future.cancel() streaming_pull_future.result()
关键改动说明
- 移除全局数据库连接,每个回调独立获取/释放连接,彻底避免多线程共享冲突
- 消息ack移到数据库写入成功之后,写入失败的消息会自动重试
- 增加全链路异常捕获,单个消息处理失败不会导致整个订阅进程崩溃
- 新增连接资源自动释放逻辑,避免连接泄漏
内容的提问来源于stack exchange,提问作者oreasech
相关产品推荐
相关产品推荐

