如何将指定SQL转换为满足累计条件的Elasticsearch DSL查询
Convert SQL Running Total Query to Elasticsearch DSL
Got it, let's break down how to translate your SQL query into Elasticsearch DSL that handles the running total logic you need. Your original SQL calculates a cumulative sum of s_b_a for each record—specifically, summing all values for the same fk_c_c_id, same month, and where datetime is less than or equal to the current record's datetime. Your current aggregation only returns monthly totals, so we need to add pipeline aggregations to track the running total over time.
Key Concepts to Replicate the SQL Logic
Your SQL's core requirements map to these Elasticsearch features:
- Group by
fk_c_c_id: Use atermsaggregation - Group by month (from
datetime): Use adate_histogramwith1Minterval - Calculate running total for records where
datetime <= current record's datetime: Use thecumulative_sumpipeline aggregation, paired with a fine-grained time histogram to capture each incremental time segment
Working Elasticsearch DSL Query
Here's the query that replicates your SQL behavior:
POST /c_c_s_index_test/_search { "size": 0, "aggs": { "group_by_customer": { "terms": { "field": "fk_c_c_id", "size": 10000 // Adjust based on your number of unique customers }, "aggs": { "group_by_month": { "date_histogram": { "field": "datetime", "interval": "1M", "format": "yyyy-MM", "min_doc_count": 1 }, "aggs": { "sort_by_time": { "date_histogram": { "field": "datetime", "interval": "1d", // Use smaller interval (e.g., 1h) if your data is more granular "min_doc_count": 1 }, "aggs": { "daily_s_b_a_sum": { "sum": { "field": "s_b_a" } }, "running_total": { "cumulative_sum": { "buckets_path": "daily_s_b_a_sum" } }, "matching_records": { "top_hits": { "size": 100, // Adjust based on how many records per time segment "_source": ["pk_c_c_s", "fk_c_c_id", "s_b_a", "datetime"] } } } } } } } } } }
How This Works
Let's walk through each part:
group_by_customer: Groups all records byfk_c_c_id, matching thei.fk_c_c_id = t.fk_c_c_idcondition in your SQL.group_by_month: Splits each customer's data into monthly buckets, matching theSUBSTR(datetime, 0,7)grouping.sort_by_time: Breaks each month into smaller time segments (days in this example) to track incremental sums. This ensures we can calculate the running total up to each point in time.daily_s_b_a_sum: Calculates the sum ofs_b_afor each daily bucket.running_total: Uses thecumulative_sumpipeline aggregation to accumulate the daily sums, giving you the total of alls_b_avalues from the start of the month up to that day—this replicates thei.datetime <= t.datetimelogic.matching_records: Usestop_hitsto return the original records associated with each daily bucket, so you can link the running total back to individual entries like your SQL does.
Notes for Adjustments
- If your
datetimevalues are more granular (e.g., hourly), change theintervalinsort_by_timeto1hinstead of1d. - Adjust the
sizeparameters intermsandtop_hitsto match your dataset's scale (too small and you'll miss data; too large and performance may suffer). - If you need the running total at the individual document level (not per day/hour), replace the
date_histograminsort_by_timewith atermsaggregation on thedatetimefield (but be cautious—this can generate a lot of buckets if your timestamps are highly unique).
内容的提问来源于stack exchange,提问作者reihaneh
相关产品推荐
相关产品推荐

