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

PySpark通过psycopg2连接PostgreSQL遇[Errno 104]连接重置求助

解决PySpark+psycopg2连接PostgreSQL超时重置问题

针对你遇到的[Errno 104] Connection reset by peer错误,以下是几个可行的解决方案:

1. 使用psycopg2连接池管理连接

单连接长时间无交互或占用过久,容易被数据库/防火墙主动断开。连接池可以自动维护连接的有效性,失效时自动重建连接。

示例代码:

from psycopg2 import pool

# 初始化连接池,根据并发需求调整minconn和maxconn
connection_pool = pool.SimpleConnectionPool(
    minconn=1,
    maxconn=5,
    database=db,
    user=user,
    password=pwd,
    host=hostname,
    port=db_port
)

# 业务处理时获取连接
con = connection_pool.getconn()
try:
    # 执行数据读写操作
    cur = con.cursor()
    # ... 你的数据处理逻辑
    con.commit()
finally:
    # 用完后归还连接到池
    connection_pool.putconn(con)

2. 开启TCP连接保活机制

通过设置psycopg2的keepalive参数,让连接定期发送心跳包,避免被判定为闲置连接而断开。

示例代码:

con = psycopg2.connect(
    database=db,
    user=user,
    password=pwd,
    host=hostname,
    port=db_port,
    keepalives=1,
    keepalives_idle=300,  # 5分钟无操作时发送心跳
    keepalives_interval=60,  # 每隔1分钟发送一次心跳
    keepalives_count=5  # 连续5次心跳失败则触发重连
)

3. 分区批量处理数据

避免单连接长时间占用,将大数据拆分为多个批次,每个批次使用独立连接(或从连接池获取)处理,处理完成后释放连接。

示例思路(PySpark分区处理):

from pyspark.sql.functions import pandas_udf
import pandas as pd

# 定义批量处理的Pandas UDF
@pandas_udf("string")
def process_batch_pandas(df: pd.DataFrame) -> pd.Series:
    conn = connection_pool.getconn()
    try:
        cur = conn.cursor()
        # 批量插入示例
        insert_query = "INSERT INTO target_table (col1, col2) VALUES (%s, %s)"
        data_tuples = list(df.itertuples(index=False, name=None))
        cur.executemany(insert_query, data_tuples)
        conn.commit()
        return pd.Series(["success"] * len(df))
    except Exception as e:
        conn.rollback()
        raise e
    finally:
        cur.close()
        connection_pool.putconn(conn)

# 按分区批量处理数据
df.mapInPandas(process_batch_pandas, schema="result string").show()

4. 调整PostgreSQL端的超时配置

检查PostgreSQL的idle_in_transaction_session_timeout参数,该参数会终止长时间处于事务空闲状态的连接。如果你的数据处理中存在大量非数据库计算(导致事务空闲),需要调整该参数:

  • 在postgresql.conf中修改:
    idle_in_transaction_session_timeout = 0  # 禁用该超时,或设置为更大的数值(如3600000毫秒=60分钟)
    
  • 修改后重启PostgreSQL,或执行ALTER SYSTEM SET idle_in_transaction_session_timeout = 0;并执行SELECT pg_reload_conf();动态生效。

内容的提问来源于stack exchange,提问作者amber_coder_backup

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 06:20:53