Apache Beam从MSSQL导数据到GCS Parquet格式报错排查
问题:Apache Beam从MSSQL读取数据写入Parquet报错TypeError
我尝试使用Apache Beam从MSSQL数据库采集数据并写入Google Cloud Storage(GCS)的Parquet文件,能成功写入CSV/TXT,但写入Parquet时触发错误:
TypeError: tuple indices must be integers or slices, not str [while running 'WriteToParquet/Write/WriteImpl/WriteBundles']
我的代码如下:
import apache_beam as beam from apache_beam.io.jdbc import ReadFromJdbc from apache_beam.typehints.schemas import LogicalType import pyarrow @LogicalType.register_logical_type class db_str(LogicalType): @classmethod def urn(cls): return "beam:logical_type:javasdk:v1" @classmethod def language_type(cls): return str def to_language_type(self, value): return str(value) def to_representation_type(self, value): return str(value) schema = pyarrow.schema([ ('CurrencyID', pyarrow.string()), ('Currency', pyarrow.string()) ]) with beam.Pipeline() as p: ip1 = (p |ReadFromJdbc( table_name='xxx.xxx', driver_class_name='com.microsoft.sqlserver.jdbc.SQLServerDriver', jdbc_url='jdbc:sqlserver://xxx.database.windows.net:1433', username='xxx', password='xxx', classpath=['com.microsoft.sqlserver:mssql-jdbc:11.2.2.jre8'], connection_properties = ';database=xxx;encrypt=true;trustServerCertificate=false;hostNameInCertificate=*.database.windows.net;loginTimeout=30;') |beam.io.WriteToParquet('gs://landingstorage/',schema=schema) )
问题原因
ReadFromJdbc默认返回元组形式的数据,而WriteToParquet要求输入是带字段名的结构化数据(如Row对象或字典)。当Parquet写入组件尝试用字符串字段名(比如'CurrencyID')索引元组时,就会触发错误——元组仅支持整数索引,不支持字符串索引。
解决方法
需要将JDBC读取到的元组转换为符合Parquet schema要求的结构化数据,以下是两种可行方案:
方案1:手动将元组转为beam.Row对象
使用beam.Map将元组映射为Row对象,确保字段名与值一一对应:
import apache_beam as beam from apache_beam.io.jdbc import ReadFromJdbc from apache_beam.typehints.schemas import Row import pyarrow schema = pyarrow.schema([ ('CurrencyID', pyarrow.string()), ('Currency', pyarrow.string()) ]) with beam.Pipeline() as p: ip1 = (p | ReadFromJdbc( table_name='xxx.xxx', driver_class_name='com.microsoft.sqlserver.jdbc.SQLServerDriver', jdbc_url='jdbc:sqlserver://xxx.database.windows.net:1433', username='xxx', password='xxx', classpath=['com.microsoft.sqlserver:mssql-jdbc:11.2.2.jre8'], connection_properties = ';database=xxx;encrypt=true;trustServerCertificate=false;hostNameInCertificate=*.database.windows.net;loginTimeout=30;' ) # 按JDBC返回的列顺序,将元组转为Row对象 | beam.Map(lambda t: Row(CurrencyID=str(t[0]), Currency=str(t[1]))) | beam.io.WriteToParquet('gs://landingstorage/', schema=schema) )
方案2:指定ReadFromJdbc的output_row_type直接返回Row对象
Apache Beam 2.40+版本支持通过output_row_type参数,让JDBC读取组件直接返回结构化的Row数据,无需手动转换:
import apache_beam as beam from apache_beam.io.jdbc import ReadFromJdbc from apache_beam.typehints.schemas import Row import pyarrow # 定义对应表结构的Row类型 class CurrencyRow(Row): CurrencyID: str Currency: str schema = pyarrow.schema([ ('CurrencyID', pyarrow.string()), ('Currency', pyarrow.string()) ]) with beam.Pipeline() as p: ip1 = (p | ReadFromJdbc( table_name='xxx.xxx', driver_class_name='com.microsoft.sqlserver.jdbc.SQLServerDriver', jdbc_url='jdbc:sqlserver://xxx.database.windows.net:1433', username='xxx', password='xxx', classpath=['com.microsoft.sqlserver:mssql-jdbc:11.2.2.jre8'], connection_properties = ';database=xxx;encrypt=true;trustServerCertificate=false;hostNameInCertificate=*.database.windows.net;loginTimeout=30;', # 指定输出的Row类型 output_row_type=CurrencyRow ) | beam.io.WriteToParquet('gs://landingstorage/', schema=schema) )
额外提示
- 你定义的
db_str逻辑类型未实际生效,因为JDBC返回的元组未关联该类型,可直接移除。 - 确保JDBC返回的列顺序与你定义的
Row字段顺序完全一致,否则会出现字段值不匹配的问题。
内容的提问来源于stack exchange,提问作者DiskoSuperStar
相关产品推荐
相关产品推荐

