使用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
相关产品推荐
相关产品推荐

