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

PyFlink应用输出[B@...而非预期CSV行,求问题排查

问题排查:PyFlink输出[B@xxx格式异常原因及修复

环境信息

  • Python 3.11.9
  • Apache Flink 2.0

需求说明

  1. 读取WeatherMeasurements.csv,列包含['city', 'state', 'county', 'avgTemp', 'avgHumidity', 'avgPressure', 'sampleTime']
  2. 解析每行数据,将华氏温度avgTemp转换为摄氏度,计算舒适指数并追加到行中
  3. 将结果写入output.csv

输入输出示例

输入(WeatherMeasurements.csv)

Madison,Wisconsin,Dane,72.2,39.94,950.89,2025-07-02T13:13:12.051886
Houston,Texas,Harris,47.03,60.53,1005.75,2025-07-02T13:14:12.051886
Madison,Wisconsin,Dane,47.26,46.11,971.65,2025-07-02T13:15:12.051886
...

预期输出(output.csv)

Madison,Wisconsin,Dane,72.20,39.94,950.89,2025-07-02T13:13:12.051886,22.33,43.36
Houston,Texas,Harris,47.03,60.53,1005.75,2025-07-02T13:14:12.051886,8.35,18.56
Madison,Wisconsin,Dane,47.26,46.11,971.65,2025-07-02T13:15:12.051886,8.48,25.47
...

实际错误输出

[B@dba59aa
[B@3be0300f
[B@7b1be48d
...

问题代码

from pyflink.datastream import StreamExecutionEnvironment, RuntimeExecutionMode
from pyflink.common import Encoder, WatermarkStrategy
from pyflink.datastream.connectors.file_system import FileSink, FileSource, StreamFormat, RollingPolicy, OutputFileConfig

# 解析每行数据为字典
def parse(line):
    city, state, county, t, h, p, ts = line.split(",")
    return {
        "city": city,
        "state": state,
        "county": county,
        "avgTemp": float(t),
        "avgHumidity": float(h),
        "avgPressure": float(p),
        "sampleTime": ts
    }

# 计算摄氏度和舒适指数并追加到字典
def enrich(rec):
    rec["avgTempC"] = (rec["avgTemp"] - 32) * 5/9
    rec["comfortIndex"] = rec["avgTemp"] * (1 - rec["avgHumidity"]/100)
    return rec

# 将字典格式化为CSV字符串
def format_csv(rec):
    csv_line = ",".join([
        rec["city"],
        rec["state"],
        rec["county"],
        f"{rec['avgTemp']:.2f}",
        f"{rec['avgHumidity']:.2f}",
        f"{rec['avgPressure']:.2f}",
        rec["sampleTime"],
        f"{rec['avgTempC']:.2f}",
        f"{rec['comfortIndex']:.2f}"
    ])
    return csv_line + "\n"

def main(file_path):
    env = StreamExecutionEnvironment.get_execution_environment()
    env.set_runtime_mode(RuntimeExecutionMode.BATCH)
    env.set_parallelism(1)

    file_src = (
        FileSource
        .for_record_stream_format(StreamFormat.text_line_format(), file_path)
        .process_static_file_set()
        .build()
    )

    lines = env.from_source(
        source=file_src,
        watermark_strategy=WatermarkStrategy.no_watermarks(),
        source_name="weather-file-source"
    )

    parsed = lines.map(parse)
    enriched = parsed.map(enrich)
    csv_lines = enriched.map(format_csv)

    csv_lines.sink_to(
        sink=FileSink.for_row_format(
            base_path="output",
            encoder=Encoder.simple_string_encoder())
        .with_output_file_config(
            OutputFileConfig.builder()
            .with_part_prefix("weather_")
            .with_part_suffix(".csv")
            .build())
        .with_rolling_policy(RollingPolicy.default_rolling_policy())
        .build()
    )

    print("Writing to output directory")
    env.execute("Weather Enrichment Job")

if __name__ == '__main__':
    file_path = "WeatherMeasurements.csv"
    main(file_path)

原因分析及修复方案

问题根源

输出的[B@xxx是Java字节数组的字符串标识,问题出在类型序列化环节:PyFlink会自动将Python字符串转换为字节数组,而Encoder.simple_string_encoder()期望接收Java字符串类型,导致最终写入的是字节数组的对象信息而非实际字符串内容。

修复方案

方案1:显式指定输出类型为Java字符串

在map(format_csv)操作时,通过output_type参数声明输出类型为Java字符串,确保PyFlink正确序列化数据:

from pyflink.common.typeinfo import Types

# ... 原有代码 ...

csv_lines = enriched.map(format_csv, output_type=Types.STRING())

方案2:改用批量写入Sink(更适合批量场景)

对于批量处理任务,使用for_bulk_format替代for_row_format,批量写入器会自动处理字符串类型转换,同时性能更优:

from pyflink.datastream.connectors.file_system import BulkWriterFactory

# ... 原有代码 ...

csv_lines.sink_to(
    sink=FileSink.for_bulk_format(
        base_path="output",
        bulk_writer_factory=BulkWriterFactory.string_bulk_writer())
    .with_output_file_config(
        OutputFileConfig.builder()
        .with_part_prefix("weather_")
        .with_part_suffix(".csv")
        .build())
    .build()
)

修复后完整代码(方案1)

from pyflink.datastream import StreamExecutionEnvironment, RuntimeExecutionMode
from pyflink.common import Encoder, WatermarkStrategy
from pyflink.datastream.connectors.file_system import FileSink, FileSource, StreamFormat, RollingPolicy, OutputFileConfig
from pyflink.common.typeinfo import Types

# 解析每行数据为字典
def parse(line):
    city, state, county, t, h, p, ts = line.split(",")
    return {
        "city": city,
        "state": state,
        "county": county,
        "avgTemp": float(t),
        "avgHumidity": float(h),
        "avgPressure": float(p),
        "sampleTime": ts
    }

# 计算摄氏度和舒适指数并追加到字典
def enrich(rec):
    rec["avgTempC"] = (rec["avgTemp"] - 32) * 5/9
    rec["comfortIndex"] = rec["avgTemp"] * (1 - rec["avgHumidity"]/100)
    return rec

# 将字典格式化为CSV字符串
def format_csv(rec):
    csv_line = ",".join([
        rec["city"],
        rec["state"],
        rec["county"],
        f"{rec['avgTemp']:.2f}",
        f"{rec['avgHumidity']:.2f}",
        f"{rec['avgPressure']:.2f}",
        rec["sampleTime"],
        f"{rec['avgTempC']:.2f}",
        f"{rec['comfortIndex']:.2f}"
    ])
    return csv_line + "\n"

def main(file_path):
    env = StreamExecutionEnvironment.get_execution_environment()
    env.set_runtime_mode(RuntimeExecutionMode.BATCH)
    env.set_parallelism(1)

    file_src = (
        FileSource
        .for_record_stream_format(StreamFormat.text_line_format(), file_path)
        .process_static_file_set()
        .build()
    )

    lines = env.from_source(
        source=file_src,
        watermark_strategy=WatermarkStrategy.no_watermarks(),
        source_name="weather-file-source"
    )

    parsed = lines.map(parse)
    enriched = parsed.map(enrich)
    # 显式指定输出类型为Java字符串
    csv_lines = enriched.map(format_csv, output_type=Types.STRING())

    csv_lines.sink_to(
        sink=FileSink.for_row_format(
            base_path="output",
            encoder=Encoder.simple_string_encoder())
        .with_output_file_config(
            OutputFileConfig.builder()
            .with_part_prefix("weather_")
            .with_part_suffix(".csv")
            .build())
        .with_rolling_policy(RollingPolicy.default_rolling_policy())
        .build()
    )

    print("Writing to output directory")
    env.execute("Weather Enrichment Job")

if __name__ == '__main__':
    file_path = "WeatherMeasurements.csv"
    main(file_path)

内容的提问来源于stack exchange,提问作者Diehardwalnut

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 18:13:12