Dask读取MariaDB遇IntCastingNaNError问题求助
问题:Dask Dataframe导出CSV触发IntCastingNaNError
使用dd.read_sql_query从MariaDB读取数据生成Dask Dataframe,尝试替换NaN/None为0后导出CSV时,报错:IntCastingNaNError: Cannot convert non-finite values (NA or inf) to integer。已尝试多种替换、填充NaN的方法,问题仍未解决。
相关代码:
#####----- #variables connection_string = 'mysql+mysqlconnector:......' indexColumn = mytitle + '.' + myrefcolumn desc_txt ='{0} desc'.format(indexColumn) #####----- #STEP 1 - SQL Query sql_query = sql.select(['*']).select_from( sql.table(mytitle) ).order_by(text(desc_txt)).limit(mylimit) #STEP 2 - read table data = dd.read_sql_query(sql_query, connection_string, index_col=myrefcolumn) #STEP 3 - replace Nan data = data.replace(to_replace=np.nan, value=0) data = data.replace(to_replace=np.inf, value=0) #data = data.replace(to_replace=np.isfinite, value=0) #--- data = data.replace(to_replace='nan', value=0) data = data.replace(to_replace='NaN', value=0) data = data.replace(to_replace='inf', value=0) data = data.replace(to_replace='NA', value=0) data = data.replace(to_replace='Na', value=0) data = data.replace(to_replace='None', value=0) data = data.replace(to_replace='NULL', value=0) data = data.replace(to_replace='NULL', value=0) data = data.replace(to_replace=' ', value=0) #data = data.replace(to_replace=None, value=0) data = data.fillna(0) #STEP 4 -- load / merger dataframes data_merge = dd.merge(dataA, dataB, on = ['idColumnA','idColumnB'], how = 'inner') #STEP 5 - exporting to csv data_merge.to_csv('csv_testing/data*.csv', index=False)
可行解决办法
方案一:SQL查询阶段直接替换NULL为0(从源头解决)
用SQLAlchemy的coalesce函数,在查询时就把数据库中的NULL值替换为0,避免后续Dask处理时出现NaN。
from sqlalchemy import func, text # 获取目标表对象 table = sql.table(mytitle) # 遍历所有字段,用coalesce将NULL替换为0,保留原字段名 select_cols = [func.coalesce(col, 0).label(col.name) for col in table.columns] # 重构查询语句 sql_query = sql.select(select_cols).select_from( table ).order_by(text(desc_txt)).limit(mylimit)
- 若为字符串类型字段,需替换为空字符串的话,把
0改成''即可。
方案二:Dask层面精准处理类型与NaN
之前的替换未生效,可能是Dask延迟执行特性或字段类型冲突导致,可按以下步骤处理:
- 先排查问题字段:
# 计算各字段的NaN数量,定位问题列 print(data.isna().sum().compute())
- 针对整数类型列单独处理:
Dask中整数列若存在NaN会自动转为float类型,需先填充再转类型:
for col in data.columns: dtype = data[col].dtype # 判定为可转为整数的float列 if dtype == 'float64' and data[col].dropna().apply(lambda x: x.is_integer()).all().compute(): data[col] = data[col].fillna(0).astype('int64') else: data[col] = data[col].fillna(0)
- 合并后再次校验处理:
合并操作可能产生新的NaN,需在导出前再处理一次:
data_merge = data_merge.fillna(0) # 重复类型校验逻辑,确保整数列无NaN for col in data_merge.columns: if data_merge[col].dtype == 'float64': if data_merge[col].dropna().apply(lambda x: x.is_integer()).all().compute(): data_merge[col] = data_merge[col].astype('int64')
方案三:过滤含NaN的行
若不需要保留含缺失值的行,可直接过滤:
- Dask层面过滤:
# 移除所有含NaN的行 data = data.dropna() data_merge = data_merge.dropna()
- SQL查询阶段过滤:
# 添加WHERE条件排除指定字段为NULL的行(按需调整字段) sql_query = sql.select(select_cols).select_from( table ).where( table.c.col1.isnot(None), table.c.col2.isnot(None) ).order_by(text(desc_txt)).limit(mylimit)
内容的提问来源于stack exchange,提问作者Brad Bailosky
相关产品推荐
相关产品推荐

