加载SQL Server 500万行数据时内核持续崩溃,Dask适配pyodbc连接报错的求助
大家好,我现在碰到个棘手的问题,想请各位帮忙看看:
我需要从SQL Server数据库里读取一张有500万行数据的大表,但不管怎么调整读取的chunk大小,Jupyter内核总是崩溃。先说明下我的运行环境(这部分是公司固定配置,我没法修改):
- 0 GPU
- 4000 CPU
- 15.0 GiB 内存
我的SQL查询代码存在项目文件夹的.sql文件里。最开始我尝试每次读取50万行的chunk,结果内核直接崩溃;改成25万行还是一样的结果;现在降到10万行,依然逃不过崩溃的命运。
按照公司规定,我必须用下面的Kerberos+pyodbc方式连接数据库,这部分代码是能正常工作的:
# Connection to SQL Server with Kerberos + pyodbc import os import pyodbc def mssql_conn_kerberos(server, driver, trusted_connection, trust_server_certificate, kerberos_cmd): # Run Kerberos for authentifications os.system(kerberos_cmd) try: # First connection attempt c_conn = pyodbc.connect( f'DRIVER={driver};' f'SERVER={server};' f'Trusted_Connection={trusted_connection};' f'TrustServerCertificate={trust_server_certificate}' ) except: # Re-run Kerberos and try authentification os.system(kerberos_cmd) c_conn = pyodbc.connect( f"DRIVER={driver};" f"SERVER={server};" f"Trusted_Connection={trusted_connection};" f"TrustServerCertificate={trust_server_certificate}" ) c_cursor = c_conn.cursor() print("Pyodbc connection ready.") return c_conn # Connection to the database
之后我写了一个函数来读取并处理.sql文件里的查询,代码如下:
import pandas as pd import time import os def call_my_query(path_to_query, query_name, chunk, connection): file_path = os.path.join(path_to_query, query_name) with open(file_path, "r") as file: query = file.read() # SQL processing in chunks + time chunks = [] start_time = time.time() for x in pd.read_sql_query(query, connection, chunksize=chunk): chunks.append(x) # Concating the chunks - joining all the chunks together df = pd.concat(chunks, ignore_index=True) # Process end-time end_time = time.time() print("Data loaded successfully!") print(f'Processed {len(df)} rows in {end_time - start_time:.2f} seconds') return df
但运行这个函数后,内核就会崩溃,提示信息如下:
The Kernel crashed while executing code in the current cell or a previous cell.
Please review the code in the cell(s) to identify a possible cause of the failure.
Click here for more info.
View Jupyter log for further details.
后来我尝试用Dask来优化读取任务,修改了call_my_query函数,但Dask和pyodbc配合时又出现了新问题。修改后的Dask版本函数如下:
import dask.dataframe as dd from sqlalchemy import select import os path_to_query = "你的查询文件路径" connection_url = "你的数据库连接URL" def call_my_query_dask(query_name, chunk, connection, index_col): # Load query from file file_path = os.path.join(path_to_query, query_name) with open(file_path, "r") as file: query_original = file.read() # Convert the SQL string/text query = select(query_original) # Start timing the process start_time = time.time() # Use Dask to read the SQL query in chunks print("Executing query and loading data with Dask...") df_dask = dd.read_sql_query( sql=query, con=connection_url, npartitions=10, index_col = index_col ) # Process end-time end_time = time.time() print("Data loaded successfully!") print(f"Processed approximately {df_dask.shape[0].compute()} rows in {end_time - start_time:.2f} seconds") return df_dask
运行这个Dask版本的函数时,会抛出以下错误:
Textual column expression 'SELECT\n\t[COL1]\n\t, [COL...' should be explicitly declared with text('SELECT\n\t[COL1]\n\t, [COL...'), or use literal_column('SELECT\n\t[COL1]\n\t, [COL...') for more specificity
实在没辙了,想请各位帮我分析下问题出在哪,有没有解决办法?谢谢大家!
备注:内容来源于stack exchange,提问作者Oskar

