Skip to content

Commit fe65f3e

Browse files
committed
refactor: Optimize connection pool #1256
1 parent 7118b40 commit fe65f3e

2 files changed

Lines changed: 34 additions & 1 deletion

File tree

‎backend/apps/datasource/crud/datasource.py‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@
2424
from ..crud.table import delete_table_by_ds_id, update_table
2525
from ..models.datasource import CoreDatasource, CreateDatasource, CoreTable, CoreField, ColumnSchema, TableObj, \
2626
DatasourceConf, TableAndFields
27+
from apps.db.db import pool_manager, driver_pool_manager
2728

2829

2930
def get_datasource_list(session: SessionDep, user: CurrentUser, oid: Optional[int] = None) -> List[CoreDatasource]:
@@ -109,6 +110,10 @@ def update_ds(session: SessionDep, trans: Trans, user: CurrentUser, ds: CoreData
109110
session.add(record)
110111
session.commit()
111112

113+
# update pool
114+
pool_manager.remove_pool(ds.id)
115+
driver_pool_manager.remove_pool(ds.id)
116+
112117
run_save_ds_embeddings([ds.id])
113118
return ds
114119

@@ -135,6 +140,11 @@ async def delete_ds(session: SessionDep, id: int):
135140
session.commit()
136141
delete_table_by_ds_id(session, id)
137142
delete_field_by_ds_id(session, id)
143+
144+
# update pool
145+
pool_manager.remove_pool(id)
146+
driver_pool_manager.remove_pool(id)
147+
138148
if term:
139149
await clear_ws_ds_cache(term.oid)
140150
return {

‎backend/apps/db/db.py‎

Lines changed: 24 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -237,7 +237,8 @@ def get_driver_connection(ds: CoreDatasource | AssistantOutDsSchema, db_config:
237237
)
238238
elif equals_ignore_case(ds.type, 'redshift'):
239239
if not use_pool:
240-
conn = redshift_connector.connect(host=conf.host, port=conf.port, database=conf.database, user=conf.username,
240+
conn = redshift_connector.connect(host=conf.host, port=conf.port, database=conf.database,
241+
user=conf.username,
241242
password=conf.password,
242243
timeout=conf.timeout, **conn_conf)
243244
else:
@@ -1175,6 +1176,17 @@ def get_pool(self, ds: CoreDatasource | AssistantOutDsSchema, **db_config):
11751176
print(f"[LRU] create: {ds.id}")
11761177
return new_pool
11771178

1179+
def remove_pool(self, datasource_id):
1180+
with self._lock:
1181+
if datasource_id in self._pools:
1182+
# 1. 从字典中移除并获取该连接池对象
1183+
pool = self._pools.pop(datasource_id)
1184+
# 2. 安全关闭该连接池,释放底层所有数据库连接和内存
1185+
pool.close()
1186+
print(f"[Manager] Closed pool and remove: {datasource_id}")
1187+
else:
1188+
print(f"[Manager] Warning: ds id {datasource_id} not exist in sqlalchemy")
1189+
11781190
def close_all(self):
11791191
"""stop"""
11801192
with self._lock:
@@ -1220,6 +1232,17 @@ def get_pool(self, ds: CoreDatasource | AssistantOutDsSchema, db_config):
12201232
print(f"[LRU] create: {ds.id}")
12211233
return new_pool
12221234

1235+
def remove_pool(self, datasource_id):
1236+
with self._lock:
1237+
if datasource_id in self._pools:
1238+
# 1. 从字典中移除并获取该连接池对象
1239+
pool = self._pools.pop(datasource_id)
1240+
# 2. 安全关闭该连接池,释放底层所有数据库连接和内存
1241+
pool.close()
1242+
print(f"[Manager] Closed pool and remove: {datasource_id}")
1243+
else:
1244+
print(f"[Manager] Warning: ds id {datasource_id} not exist in dbutils")
1245+
12231246
def close_all(self):
12241247
"""stop"""
12251248
with self._lock:

0 commit comments

Comments
 (0)