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

能否用Java Arrow实现PyArrow的Parquet读写与Hive分区写入功能?

使用Apache Arrow Java实现等效功能

原PyArrow代码核心功能为:读取本地Parquet文件,按指定列以Hive风格分区写入数据集,写入时删除匹配的现有分区数据,同时启用多线程并限制最大分区数。以下是Java环境的实现方案:

1. 依赖配置(Maven)

需引入Apache Arrow核心、数据集及Parquet相关依赖:

<dependencies>
    <dependency>
        <groupId>org.apache.arrow</groupId>
        <artifactId>arrow-core</artifactId>
        <version>15.0.0</version> <!-- 使用最新稳定版 -->
    </dependency>
    <dependency>
        <groupId>org.apache.arrow</groupId>
        <artifactId>arrow-dataset</artifactId>
        <version>15.0.0</version>
    </dependency>
    <dependency>
        <groupId>org.apache.arrow</groupId>
        <artifactId>arrow-parquet</artifactId>
        <version>15.0.0</version>
    </dependency>
    <dependency>
        <groupId>org.apache.arrow</groupId>
        <artifactId>arrow-memory-netty</artifactId>
        <version>15.0.0</version>
    </dependency>
</dependencies>

2. Apache Arrow Java代码实现

import org.apache.arrow.dataset.file.FileDatasetFactory;
import org.apache.arrow.dataset.file.FileFormat;
import org.apache.arrow.dataset.jni.NativeMemoryPool;
import org.apache.arrow.dataset.scanner.ScanOptions;
import org.apache.arrow.dataset.scanner.Scanner;
import org.apache.arrow.dataset.source.Dataset;
import org.apache.arrow.dataset.write.DatasetWriter;
import org.apache.arrow.dataset.write.DatasetWriterOptions;
import org.apache.arrow.dataset.write.WriteConfig;
import org.apache.arrow.memory.BufferAllocator;
import org.apache.arrow.memory.RootAllocator;

import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.Paths;

public class ParquetHivePartitionWriter {
    public static void main(String[] args) {
        String fileName = args[0];
        String inputPath = "/Users/some-user/Downloads/" + fileName + ".parquet";
        String outputBaseDir = "/Users/some-user/hive";
        String partitionColumn = "column";
        int maxPartitions = 10000;

        try (BufferAllocator allocator = new RootAllocator();
             NativeMemoryPool memoryPool = NativeMemoryPool.getDefault()) {

            // 读取Parquet文件为Arrow数据集
            Path inputFile = Paths.get(inputPath);
            FileDatasetFactory factory = new FileDatasetFactory(
                    allocator,
                    memoryPool,
                    FileFormat.PARQUET,
                    inputFile.toUri().toString()
            );
            Dataset dataset = factory.create();

            // 配置写入参数
            WriteConfig writeConfig = WriteConfig.builder()
                    .withMaxPartitions(maxPartitions)
                    .withUseThreads(true) // 启用多线程写入
                    .build();

            DatasetWriterOptions writerOptions = DatasetWriterOptions.builder()
                    .fileFormat(FileFormat.PARQUET)
                    .partitionColumns(partitionColumn)
                    .partitionType(DatasetWriterOptions.PartitionType.HIVE) // Hive风格分区目录
                    .writeConfig(writeConfig)
                    .build();

            // 实现existing_data_behavior='delete_matching':删除现有匹配分区
            Path outputDir = Paths.get(outputBaseDir);
            if (Files.exists(outputDir)) {
                // 遍历并删除以目标分区列命名的目录
                Files.walk(outputDir)
                        .filter(Files::isDirectory)
                        .filter(path -> path.getFileName().toString().startsWith(partitionColumn + "="))
                        .forEach(path -> {
                            try {
                                Files.delete(path);
                            } catch (Exception e) {
                                e.printStackTrace();
                            }
                        });
            }

            // 写入分区数据集
            try (DatasetWriter writer = DatasetWriter.create(
                    outputDir.toUri().toString(),
                    writerOptions,
                    allocator,
                    memoryPool
            )) {
                Scanner scanner = dataset.newScan(ScanOptions.builder().build());
                writer.write(scanner);
            }

        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

3. 替代方案:Apache Spark Java API

若需更简洁的Hive分区管理,可使用Spark API,其内置分区覆盖逻辑,无需手动处理删除:

import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SparkSession;
import org.apache.spark.sql.SaveMode;

public class SparkParquetHiveWriter {
    public static void main(String[] args) {
        SparkSession spark = SparkSession.builder()
                .appName("ParquetHivePartitionWriter")
                .master("local[*]") // 生产环境移除该配置
                .getOrCreate();

        String fileName = args[0];
        String inputPath = "/Users/some-user/Downloads/" + fileName + ".parquet";
        String outputBaseDir = "/Users/some-user/hive";

        // 读取Parquet文件
        Dataset<Row> df = spark.read().parquet(inputPath);

        // 写入Hive分区,SaveMode.Overwrite自动覆盖匹配分区
        df.write()
                .mode(SaveMode.Overwrite)
                .partitionBy("column")
                .parquet(outputBaseDir);

        spark.stop();
    }
}

关键说明

  • Arrow方案:需手动实现匹配分区删除逻辑,适合轻量、依赖Arrow生态的场景
  • Spark方案:SaveMode.Overwrite直接对应existing_data_behavior='delete_matching',多线程、分区管理更省心,适合大数据场景

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 17:35:26