如何高效将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
相关产品推荐
相关产品推荐

