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

如何高效将ElasticSearch响应转为Parquet文件(动态生成Schema)

解决方案

1. 避开Spark,用Parquet原生API+Jackson处理(低内存占用)

Spark更适合分布式大数据场景,当内存资源受限的时候,直接用Parquet Java API配合Jackson解析ES响应,内存占用会更低、效率更高。

核心流程:

  • 从SearchResponse<ObjectNode>中提取所有命中数据的_source(response.hits().hits().stream().map(h -> h.source()))
  • 基于第一条ObjectNode递归生成支持嵌套结构的Parquet Schema
  • 逐条将ObjectNode转换为Parquet记录写入文件

代码实现

首先添加Maven依赖:

<dependencies>
    <dependency>
        <groupId>org.apache.parquet</groupId>
        <artifactId>parquet-hadoop</artifactId>
        <version>1.14.0</version>
    </dependency>
    <dependency>
        <groupId>com.fasterxml.jackson.core</groupId>
        <artifactId>jackson-databind</artifactId>
        <version>2.15.2</version>
    </dependency>
    <dependency>
        <groupId>org.apache.hadoop</groupId>
        <artifactId>hadoop-common</artifactId>
        <version>3.3.6</version>
        <scope>provided</scope>
    </dependency>
</dependencies>

转换工具类代码:

import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ObjectNode;
import co.elastic.clients.elasticsearch.core.SearchResponse;
import co.elastic.clients.elasticsearch.core.hits.Hit;
import org.apache.parquet.hadoop.ParquetWriter;
import org.apache.parquet.hadoop.api.WriteSupport;
import org.apache.parquet.io.api.Binary;
import org.apache.parquet.io.api.RecordConsumer;
import org.apache.parquet.schema.*;

import java.io.File;
import java.io.IOException;
import java.util.Iterator;
import java.util.Map;

public class EsToParquetConverter {

    private static final ObjectMapper MAPPER = new ObjectMapper();

    public static void convert(SearchResponse<ObjectNode> esResponse, File outputParquet) throws IOException {
        Iterator<Hit<ObjectNode>> hitIterator = esResponse.hits().hits().iterator();
        if (!hitIterator.hasNext()) {
            return;
        }

        // 用第一条数据生成Parquet Schema
        ObjectNode firstSource = hitIterator.next().source();
        MessageType parquetSchema = buildParquetSchema(firstSource, "es_data");

        // 初始化Parquet写入器
        try (ParquetWriter<JsonNode> writer = new ParquetWriter<>(
                outputParquet,
                new JsonNodeWriteSupport(parquetSchema)
        )) {
            writer.write(firstSource);
            // 批量写入剩余数据
            hitIterator.forEachRemaining(hit -> {
                try {
                    writer.write(hit.source());
                } catch (IOException e) {
                    throw new RuntimeException("写入Parquet文件失败", e);
                }
            });
        }
    }

    // 递归构建支持嵌套结构的Parquet Schema
    private static MessageType buildParquetSchema(ObjectNode rootNode, String schemaName) {
        Types.MessageTypeBuilder builder = Types.buildMessage();
        for (Map.Entry<String, JsonNode> field : rootNode.fields()) {
            builder.addField(buildFieldType(field.getKey(), field.getValue()));
        }
        return builder.named(schemaName);
    }

    private static Type buildFieldType(String fieldName, JsonNode node) {
        if (node.isObject()) {
            Types.GroupBuilder groupBuilder = Types.requiredGroup();
            for (Map.Entry<String, JsonNode> nestedField : node.fields()) {
                groupBuilder.addField(buildFieldType(nestedField.getKey(), nestedField.getValue()));
            }
            return groupBuilder.named(fieldName);
        } else if (node.isArray()) {
            if (node.size() > 0) {
                Type elementType = buildFieldType("element", node.get(0));
                return Types.requiredList(elementType).named(fieldName);
            } else {
                return Types.optionalList(Types.optionalPrimitive(PrimitiveType.PrimitiveTypeName.BINARY)).named(fieldName);
            }
        } else if (node.isTextual()) {
            return Types.requiredBinary().named(fieldName);
        } else if (node.isInt()) {
            return Types.requiredInt32().named(fieldName);
        } else if (node.isLong()) {
            return Types.requiredInt64().named(fieldName);
        } else if (node.isFloat()) {
            return Types.requiredFloat().named(fieldName);
        } else if (node.isDouble()) {
            return Types.requiredDouble().named(fieldName);
        } else if (node.isBoolean()) {
            return Types.requiredBoolean().named(fieldName);
        } else {
            return Types.requiredBinary().named(fieldName);
        }
    }

    // 自定义WriteSupport,实现JsonNode到Parquet记录的转换
    private static class JsonNodeWriteSupport extends WriteSupport<JsonNode> {
        private final MessageType schema;
        private RecordConsumer consumer;

        public JsonNodeWriteSupport(MessageType schema) {
            this.schema = schema;
        }

        @Override
        public WriteContext init(Map<String, String> configuration) {
            return new WriteContext(schema, configuration);
        }

        @Override
        public void prepareForWrite(RecordConsumer recordConsumer) {
            this.consumer = recordConsumer;
        }

        @Override
        public void write(JsonNode record) {
            consumer.startMessage();
            writeNode(record, schema.getFields().iterator());
            consumer.endMessage();
        }

        private void writeNode(JsonNode node, Iterator<Type> fieldIterator) {
            while (fieldIterator.hasNext()) {
                Type field = fieldIterator.next();
                String fieldName = field.getName();
                JsonNode fieldNode = node.get(fieldName);
                if (fieldNode == null || fieldNode.isNull()) {
                    continue;
                }

                consumer.startField(fieldName, 0);
                if (field.isPrimitive()) {
                    PrimitiveType primitiveType = (PrimitiveType) field;
                    switch (primitiveType.getPrimitiveTypeName()) {
                        case BINARY:
                            consumer.addBinary(Binary.fromString(fieldNode.asText()));
                            break;
                        case INT32:
                            consumer.addInteger(fieldNode.asInt());
                            break;
                        case INT64:
                            consumer.addLong(fieldNode.asLong());
                            break;
                        case FLOAT:
                            consumer.addFloat(fieldNode.floatValue());
                            break;
                        case DOUBLE:
                            consumer.addDouble(fieldNode.doubleValue());
                            break;
                        case BOOLEAN:
                            consumer.addBoolean(fieldNode.asBoolean());
                            break;
                        default:
                            consumer.addBinary(Binary.fromString(fieldNode.toString()));
                    }
                } else if (field.isGroup()) {
                    GroupType groupType = (GroupType) field;
                    consumer.startGroup();
                    writeNode(fieldNode, groupType.getFields().iterator());
                    consumer.endGroup();
                } else if (field.isList()) {
                    GroupType listType = (GroupType) field;
                    consumer.startGroup();
                    for (JsonNode arrayElement : fieldNode) {
                        consumer.startField("element", 0);
                        if (listType.getFields().get(0).isPrimitive()) {
                            PrimitiveType elemType = (PrimitiveType) listType.getFields().get(0);
                            switch (elemType.getPrimitiveTypeName()) {
                                case BINARY:
                                    consumer.addBinary(Binary.fromString(arrayElement.asText()));
                                    break;
                                case INT32:
                                    consumer.addInteger(arrayElement.asInt());
                                    break;
                                case INT64:
                                    consumer.addLong(arrayElement.asLong());
                                    break;
                            }
                        } else {
                            consumer.startGroup();
                            writeNode(arrayElement, listType.getFields().get(0).asGroupType().getFields().iterator());
                            consumer.endGroup();
                        }
                        consumer.endField("element", 0);
                    }
                    consumer.endGroup();
                }
                consumer.endField(fieldName, 0);
            }
        }
    }
}

在异步滚动方法中调用:

asyncClient.search(searchRequest, ObjectNode.class).whenCompleteAsync((response, ex) -> {
    if (ex != null) {
        promise.fail(ex);
        return;
    }
    try {
        File outputFile = new File("async.parquet");
        EsToParquetConverter.convert(response, outputFile);
        promise.complete(new JsonObject().put("status", "success"));
    } catch (IOException e) {
        promise.fail(e);
    }
});

2. 优化Spark方式(若必须使用Spark)

如果一定要用Spark,可通过以下方式解决内存溢出问题:

  • 不要直接读取response.toString(),而是提取hits中的_source分批构建RDD/DataFrame
  • 调整Spark堆内存参数(--driver-memory和--executor-memory)
  • 使用Spark官方ES连接器直接读取数据,自动处理分页和Schema推断:
Dataset<Row> df = sparkSession.read()
        .format("org.elasticsearch.spark.sql")
        .option("es.nodes", "你的ES地址")
        .option("es.port", "9200")
        .option("es.query", query.toString())
        .load("目标索引名");
df.write().parquet("async.parquet");

3. 关于嵌套结构Schema生成的疑问

完全支持嵌套结构。JSON的嵌套对象/数组可直接映射为Parquet的GroupType和ListType,通过递归解析JSON结构即可生成对应Schema。需要注意:

  • 若后续记录存在新增字段,首条记录生成的Schema会缺失这些字段,可收集多条记录合并Schema,或预设字段为可选类型
  • 尽量保证数组内元素类型一致,否则需做兼容处理(如统一存储为字符串)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 09:59:52