使用Dask读取PostgreSQL表时触发TypeError问题求助
使用Dask读取PostgreSQL表时出现TypeError的解决方法
问题代码
import dask.dataframe as dd df_temporal = dd.read_sql_table('daily_data', 'postgresql://schmid@cedar-pgsql-vm:5555/schmid_db', 'id') print("Number null values: ", df_temporal['waterbody_id'].isnull().count().compute())
报错信息
Traceback (most recent call last): File "/project/6005341/schmidj/Preprocessing_temporal_data.py", line 46, in <module> print("Number null values: ", df_temporal['waterbody_id'].isnull().count().compute()) File "/home/schmidj/ANGLERENV/lib/python3.10/site-packages/dask/base.py", line 315, in compute (result,) = compute(self, traverse=False, **kwargs) File "/home/schmidj/ANGLERENV/lib/python3.10/site-packages/dask/base.py", line 600, in compute results = schedule(dsk, keys, **kwargs) File "/home/schmidj/ANGLERENV/lib/python3.10/site-packages/dask/threaded.py", line 89, in get results = get_async( File "/home/schmidj/ANGLERENV/lib/python3.10/site-packages/dask/local.py", line 511, in get_async raise_exception(exc, tb) File "/home/schmidj/ANGLERENV/lib/python3.10/site-packages/dask/local.py", line 319, in reraise raise exc File "/home/schmidj/ANGLERENV/lib/python3.10/site-packages/dask/local.py", line 224, in execute_task result = _execute_task(task, data) File "/home/schmidj/ANGLERENV/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task return func(*(_execute_task(a, cache) for a in args)) File "/home/schmidj/ANGLERENV/lib/python3.10/site-packages/dask/core.py", line 119, in <genexpr> return func(*(_execute_task(a, cache) for a in args)) File "/home/schmidj/ANGLERENV/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task return func(*(_execute_task(a, cache) for a in args)) File "/home/schmidj/ANGLERENV/lib/python3.10/site-packages/dask/utils.py", line 71, in apply return func(*args, **kwargs) File "/home/schmidj/ANGLERENV/lib/python3.10/site-packages/dask/dataframe/io/sql.py", line 432, in _read_sql_chunk return df.astype(meta.dtypes.to_dict(), copy=False) File "/home/schmidj/ANGLERENV/lib/python3.10/site-packages/pandas/core/generic.py", line 6231, in astype res_col = col.astype(dtype=cdt, copy=copy, errors=errors) File "/home/schmidj/ANGLERENV/lib/python3.10/site-packages/pandas/core/generic.py", line 6245, in astype new_data = self._mgr.astype(dtype=dtype, copy=copy, errors=errors) File "/home/schmidj/ANGLERENV/lib/python3.10/site-packages/pandas/core/internals/managers.py", line 446, in astype return self.apply("astype", dtype=dtype, copy=copy, errors=errors) File "/home/schmidj/ANGLERENV/lib/python3.10/site-packages/pandas/core/internals/managers.py", line 348, in apply applied = getattr(b, f)(**kwargs) File "/home/schmidj/ANGLERENV/lib/python3.10/site-packages/pandas/core/internals/blocks.py", line 527, in astype new_values = astype_array_safe(values, dtype, copy=copy, errors=errors) File "/home/schmidj/ANGLERENV/lib/python3.10/site-packages/pandas/core/dtypes/astype.py", line 299, in astype_array_safe new_values = astype_array(values, dtype, copy=copy) File "/home/schmidj/ANGLERENV/lib/python3.10/site-packages/pandas/core/dtypes/astype.py", line 230, in astype_array values = astype_nansafe(values, dtype, copy=copy) File "/home/schmidj/ANGLERENV/lib/python3.10/site-packages/pandas/core/dtypes/astype.py", line 170, in astype_nansafe return arr.astype(dtype, copy=True) TypeError: int() argument must be a string, a bytes-like object or a real number, not 'NoneType'
排查情况
- 同一数据库下其他表(如
waterbody表)可正常读取,问题集中在daily_data表 - 查询表结构的SQL语句:
SELECT column_name, data_type, character_maximum_length, table_name,ordinal_position, is_nullable FROM information_schema.COLUMNS WHERE table_name LIKE 'daily_data' ORDER BY ordinal_position;
- 结果显示
waterbody_id为不可为空的integer类型,但Dask读取后,多个应为数值类型的字段(如total_trips、hours_out)被识别为object类型
解决方法
1. 手动指定元数据(Meta)
Dask自动推断元数据时可能出错,手动指定各字段的正确类型可以避免类型转换错误:
import dask.dataframe as dd import pandas as pd # 定义元数据,根据实际表结构调整字段和类型 meta = pd.DataFrame({ 'id': pd.Int64Dtype(), 'waterbody_id': pd.Int64Dtype(), 'total_trips': pd.Int64Dtype(), 'hours_out': pd.Float64Dtype(), # 添加其他字段及其对应类型 }) df_temporal = dd.read_sql_table( 'daily_data', 'postgresql://schmid@cedar-pgsql-vm:5555/schmid_db', 'id', meta=meta )
2. 检查表中实际数据
虽然表结构标记waterbody_id不可为空,但可能存在脏数据(比如实际存储了NULL或者非数值内容),可以用SQL检查:
-- 检查waterbody_id是否存在NULL或非数值 SELECT * FROM daily_data WHERE waterbody_id IS NULL OR waterbody_id !~ '^[0-9]+$';
如果存在脏数据,需要先清理数据库中的数据,再进行读取。
3. 调整读取时的类型转换策略
在read_sql_table中添加dtype参数,指定字段的类型:
df_temporal = dd.read_sql_table( 'daily_data', 'postgresql://schmid@cedar-pgsql-vm:5555/schmid_db', 'id', dtype={'waterbody_id': 'Int64', 'total_trips': 'Int64', 'hours_out': 'float64'} )
4. 先读取小样本验证类型
先读取部分数据,确认类型后再读取全量:
# 先读取前100行 sample_df = dd.read_sql_table( 'daily_data', 'postgresql://schmid@cedar-pgsql-vm:5555/schmid_db', 'id', npartitions=1 ).head(100) # 查看样本的类型,调整元数据 print(sample_df.dtypes)
内容的提问来源于stack exchange,提问作者Julia2203
相关产品推荐
相关产品推荐

