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

使用Akka写入Parquet Sink异常:数据未写入流程已结束

解决Akka流写入Parquet未完成就结束的问题

核心问题分析

  • FileIO读取的是字节块而非CSV行:你当前用FileIO.fromPath直接读取字节块后拆分,会导致跨行内容被错误切割,甚至出现半行数据,后续构建GenericRecord时极易抛出异常,导致流静默终止(Akka流默认异常时终止流,未处理则无日志输出)。
  • ParquetWriter生命周期管理错误:在流完成回调中手动关闭writer时,Sink可能仍在异步写入数据,提前关闭会导致数据丢失;同时自定义Sink若未正确处理写入完成逻辑,会导致流提前结束。
  • 缺乏异常监控:流中转换异常(如CSV解析错误、Integer转换失败)未被捕获日志,无法定位问题根源。

修正方案

1. 引入CSV解析依赖

使用Akka官方CSV解析库按行处理CSV,避免字节块拆分问题(以Maven为例):

<dependency>
    <groupId>com.typesafe.akka</groupId>
    <artifactId>akka-stream-csv_2.13</artifactId>
    <version>2.6.20</version>
</dependency>

2. 修正流处理逻辑

重新构建流,确保按行解析CSV、正确管理ParquetWriter生命周期:

String parquetOut = "out.parquet";
String schemaFilePath = "schema.json";
final String path = "data.csv";
final java.nio.file.Path file = Paths.get(path);
ActorSystem<Void> actorSystem = ActorSystem.create(Behaviors.empty(), "actorSystem");

// Avro Parquet 初始化
Configuration conf = new Configuration();
conf.setBoolean(AvroReadSupport.AVRO_COMPATIBILITY, true);
Schema schema = getSchema(schemaFilePath);

// 用Akka资源管理自动维护ParquetWriter生命周期
CompletionStage<IOResult> result = Source.fromResource(() -> 
        AvroParquetWriter.<GenericRecord>builder(new Path(parquetOut))
                .withConf(conf)
                .withWriteMode(ParquetFileWriter.Mode.OVERWRITE)
                .withCompressionCodec(CompressionCodecName.SNAPPY)
                .withDictionaryEncoding(true)
                .withSchema(schema)
                .build()
    )
    .flatMapConcat(writer -> 
        FileIO.fromPath(file)
            // 按行解析CSV字节流
            .via(CsvParsing.lineScanner())
            // 将每行字节转换为字符串列表
            .map(row -> row.map(ByteString::utf8String).toList())
            // 转换为GenericRecord并添加日志验证
            .map(values -> {
                GenericRecord record = new GenericData.Record(schema);
                record.put("first-name", values.get(0));
                record.put("last-name", values.get(1));
                record.put("id", Integer.valueOf(values.get(2)));
                record.put("age", Integer.valueOf(values.get(3)));
                System.out.println("已构建Record: " + record);
                return record;
            })
            // 监控流中异常
            .log("parquet-process-error")
            .addAttributes(Attributes.logLevels(
                LoggingLevel.ERROR, LoggingLevel.ERROR, LoggingLevel.ERROR))
            // 写入Parquet并确保处理完成
            .toMat(Sink.foreach(record -> {
                try {
                    writer.write(record);
                } catch (IOException e) {
                    throw new RuntimeException(e);
                }
            }), Keep.left())
    )
    .run(actorSystem);

result.whenComplete((input, exception) -> {
    if (exception != null) {
        System.err.println("流执行异常: " + exception.getMessage());
        exception.printStackTrace();
    } else {
        System.out.println("流执行完成,处理结果: " + input);
    }
    actorSystem.terminate();
});

3. 关键改进点

  • 按行解析CSV:用CsvParsing.lineScanner()将字节流正确拆分为CSV行,避免半行/跨行问题。
  • 自动资源管理:Source.fromResource会在流结束后自动关闭ParquetWriter,无需手动调用close()。
  • 异常监控:通过.log()操作符捕获流中异常,便于排查问题。
  • 异常包装:将IO异常转为RuntimeException,确保Akka流能捕获并终止流程。

额外排查建议

  • 检查CSV文件是否存在格式错误(空行、字段数不匹配),避免values.get(index)抛出下标越界异常。
  • 确认Avro Schema的字段名、类型与CSV数据完全匹配,比如id/age是否为int类型。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.23 05:15:40