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
相关产品推荐
相关产品推荐

