如何将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
- Replace
map()withstruct(): Instead of creating key-value pairs,struct()builds structured objects with explicit field names that align with your target JSON structure. - Collect structs into a list: Use
collect_list()to gather these structured objects into an array grouped by theproductfield. - 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
相关产品推荐
相关产品推荐

