PyFlink应用输出[B@...而非预期CSV行,求问题排查
问题排查:PyFlink输出
[B@xxx格式异常原因及修复 环境信息
- Python 3.11.9
- Apache Flink 2.0
需求说明
- 读取
WeatherMeasurements.csv,列包含['city', 'state', 'county', 'avgTemp', 'avgHumidity', 'avgPressure', 'sampleTime'] - 解析每行数据,将华氏温度
avgTemp转换为摄氏度,计算舒适指数并追加到行中 - 将结果写入
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
相关产品推荐
相关产品推荐

