使用Apache Beam向PostgreSQL写入数据时类型转换失败求助
Apache Beam写入PostgreSQL时字符串列类型转换失败问题
问题描述
使用Apache Beam读取CSV文件并写入PostgreSQL时,管道因类型转换不匹配失败。当两列均为int类型时管道可正常运行,但包含字符串类型的列时就会失败。
已尝试的方案
方案1:用NamedTuple指定字段类型
from past.builtins import unicode import typing import apache_beam as beam from apache_beam import convert from apache_beam.io.jdbc import WriteToJdbc ExampleRow = typing.NamedTuple('ExampleRow',[('id',int),('name',unicode)]) beam_df = (pipeline | 'Read CSV' >> beam.dataframe.io.read_csv('path.csv').with_output_types(ExampleRow)) beam_df2 = (convert.to_pcollection(beam_df) | beam.Map(print) | WriteToJdbc( table_name=table_name, jdbc_url=jdbc_url, driver_class_name = 'org.postgresql.Driver', statement="insert into tablr values(?,?);", username=username, password=password, ) ) result = pipeline.run() result.wait_until_finish()
方案2:自定义LogicalType转换字符串
from apache_beam.typehints.schemas import LogicalType from past.builtins import unicode @LogicalType.register_logical_type class db_str(LogicalType): @classmethod def urn(cls): return "beam:logical_type:javasdk:v1" @classmethod def language_type(cls): return unicode def to_language_type(self, value): return unicode(value) def to_representation_type(self, value): return unicode(value)
补充信息
打印的PCollection内容如下:
BeamSchema_f0d95d64_95c7_43ba_8a04_ac6a0b7352d9(id=21, nom='nom21') BeamSchema_f0d95d64_95c7_43ba_8a04_ac6a0b7352d9(id=22, nom='nom22') BeamSchema_f0d95d64_95c7_43ba_8a04_ac6a0b7352d9(id=21, nom='nom21') BeamSchema_f0d95d64_95c7_43ba_8a04_ac6a0b7352d9(id=22, nom='nom22')
问题明确出在WriteToJdbc函数和nom列的类型转换上,求解决方法。
解决方法
1. 修正JDBC插入语句的字段映射
首先检查插入语句:你代码里写的是insert into tablr values(?,?);,但打印的字段是id和nom,要确保PostgreSQL表的字段顺序、名称和类型完全匹配,表中nom字段需设为varchar或text类型,建议显式指定字段名避免顺序错误:
insert into tablr(id, nom) values(?, ?);
2. 替换unicode为原生str类型
Python 3中已移除unicode类型,改用标准str定义NamedTuple,避免兼容性问题:
# 去掉from past.builtins import unicode ExampleRow = typing.NamedTuple('ExampleRow',[('id', int), ('name', str)])
3. 手动解析CSV,跳过DataFrame类型推断
绕过DataFrame的隐式类型转换问题,直接用ReadFromText读取并手动解析每行数据,确保类型明确:
import csv from io import StringIO import apache_beam as beam from apache_beam.io.jdbc import WriteToJdbc def parse_csv_line(line): # 按CSV列顺序解析,第一列转int,第二列保留str reader = csv.reader(StringIO(line)) for row in reader: return (int(row[0]), row[1]) with beam.Pipeline() as pipeline: (pipeline | 'Read CSV' >> beam.io.ReadFromText('path.csv', skip_header_lines=1) | 'Parse Lines' >> beam.Map(parse_csv_line) | 'Write to PostgreSQL' >> WriteToJdbc( table_name=table_name, jdbc_url=jdbc_url, driver_class_name='org.postgresql.Driver', statement="insert into tablr(id, nom) values(?, ?);", username=username, password=password ) )
4. 检查JDBC驱动兼容性
确保使用的PostgreSQL JDBC驱动版本与Apache Beam版本匹配,建议安装最新稳定版的驱动:
pip install psycopg2-binary
内容的提问来源于stack exchange,提问作者Ghassen Sultana
相关产品推荐
相关产品推荐

