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

如何将Spark数据集列聚合为JSON数组以适配ElasticSearch存储?

Solution to Convert Spark Aggregation to Target JSON Format

Got it, let's tweak your Spark code to get exactly the JSON format you need! The problem with your current approach is that using map() gives you key-value pairs (like Roger -> 36), but you need structured objects with clear field names (purchaser and amount). Here's how to fix this step by step:

Step-by-Step Breakdown

  1. Replace map() with struct(): Instead of creating key-value pairs, struct() builds structured objects with explicit field names that align with your target JSON structure.
  2. Collect structs into a list: Use collect_list() to gather these structured objects into an array grouped by the product field.
  3. Convert to JSON string: Use to_json() to turn the array of structs into a properly formatted JSON string, and rename the column to match your desired output.

Modified Code

import org.apache.spark.sql.functions.{struct, collect_list, to_json, col};

// Your original data setup stays the same
List<Bean> data = new ArrayList<Bean>();
data.add(new Bean("book","John",59));
data.add(new Bean("book","Björn",61));
data.add(new Bean("tv","Roger",36));
Dataset ds = spark.createDataFrame(data, Bean.class);

// Updated aggregation logic
ds = ds.groupBy(col("product"))
       .agg(collect_list(struct(col("purchaser"), col("amount"))).alias("purchases"))
       .withColumn("json", to_json(col("purchases")))
       .drop("purchases"); // Optional: remove the intermediate "purchases" column

// Show the final result
ds.show(false);

Final Output

+-------+------------------------------------------------------------------+
|product|json                                                              |
+-------+------------------------------------------------------------------+
|tv     |[{"purchaser":"Roger","amount":36}]                                |
|book   |[{"purchaser":"John","amount":59},{"purchaser":"Björn","amount":61}]|
+-------+------------------------------------------------------------------+

Key Details

  • struct(col("purchaser"), col("amount")) generates structured objects that directly translate to JSON objects like {"purchaser":"Roger","amount":36}.
  • collect_list(...) aggregates these objects into an array for each unique product.
  • to_json() handles all the JSON formatting automatically, ensuring the output matches your required structure perfectly.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 10:09:10