You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.05 23:27:03