You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.03 19:15:37