【发布时间】:2022-01-27 11:46:17
【问题描述】:
在下面的示例代码中,在我们的Azure Databricks 笔记本的one cell 中,代码将大约2000 万条记录从Azure SQL db 加载到Python pandas dataframe,通过应用一些函数进行一些数据帧列转换(如下面的代码sn-p所示)。但是在运行代码大约半小时后,Databricks 抛出以下错误:
错误:
ConnectException: Connection refused (Connection refused)
Error while obtaining a new communication channel
ConnectException error: This is often caused by an OOM error that causes the connection to the Python REPL to be closed. Check your query's memory usage.
备注:表格大约有 150 列。 Databricks 上的Spark setting 如下:
集群:128 GB , 16 Cores, DBR 8.3, Spark 8.3, Scala 2.12
问题:错误的原因可能是什么,我们该如何解决?
import sqlalchemy as sq
import pandas as pd
def fn_myFunction(lastname):
testvar = lastname.lower()
testvar = testvar.strip()
return testvar
pw = dbutils.secrets.get(scope='SomeScope',key='sql')
engine = sq.create_engine('mssql+pymssql://SERVICE.Databricks.NONPUBLICETL:'+pw+'MyAzureSQL.database.windows.net:1433/TEST', isolation_level="AUTOCOMMIT")
app_df = pd.read_sql('select * from MyTable', con=engine)
#create new column
app_df['NewColumn'] = app_df['TestColumn'].apply(lambda x: fn_myFunction(x))
.............
.............
【问题讨论】:
标签: python pandas apache-spark azure-sql-database azure-databricks