My documents in ElasticSearch look like the following:
{
"_id" : "Eba714EBOZm74KGk5PJr",
"_source" : {
"customer-id" : "1",
"customer-balance" : 15,
"eventDate" : 1640995200000 //01-jan-22
}
},
{
"_id" : "z7a814EBOZm74KGkHfLv",
"_source" : {
"customer-id" : "1",
"customer-balance" : 20,
"eventDate" : 1641081600000 //02-jan-22
}
},
{
"_id" : "dba814EBOZm74KGkNfPE",
"_source" : {
"customer-id" : "1",
"customer-balance" : 25,
"eventDate" : 1641168000000 //03-jan-22
}
},
{
"_id" : "5Le814EBOZm74KGkcQO_",
"_source" : {
"customer-id" : "2",
"customer-balance" : 15,
"eventDate" : 1640995200000 //01-jan-22
}
}
My data is event-driven, I update elastic search when it changes (i.e. when I get new data). This data however comes in at different cadences, some customers may be updated by the minute and others weekly. I want to bucket data by arbitrary periods (e.g. daily) and on search aggregate the total balance across customers for each bucket. To do this I can write a query like the following:
GET my_index/_search
{
"aggs": {
"balance_over_time": {
"date_histogram": {
"field": "eventDate",
"fixed_interval": "1d"
},
"aggs": {
"daily-balances": {
"sum": {
"field": "customer-balance"
}
}
}
}
}
}
We get the following result:
"aggregations" : {
"balance_over_time" : {
"buckets" : [
{
"key" : 1640995200000,
"doc_count" : 2,
"daily-balances" : {
"value" : 30.0
}
},
{
"key" : 1641081600000,
"doc_count" : 1,
"daily-balances" : {
"value" : 20.0
}
},
{
"key" : 1641168000000,
"doc_count" : 1,
"daily-balances" : {
"value" : 25.0
}
}
]
}
You will note that customer "2" hasn't been updated on days 2 or 3 - I would like their balance to therefore roll forward so that I get daily balances of 30, 35, 40 instead of 30, 20, 25.
How can I do this? A few other limitations:
- It isn't just 1 balance field, I have dozens of these fields.
- Ideally I want to allow my queries to not be restricted in the bucket intervals, e.g. I may want to use 5 mins intervals or may want to use 30-day intervals and would like to decide on the fly
- I also don't want to include the same customer twice if they have two updates in a period - ideally I could choose the first/last/average update wins approach.
Is any of this viable? Am I trying to use Elastic to do something it's just not good at?
My elastic is being hosted by elastic.co and I'm using logstash to create the data, sourced from a Kafka topic, if any of that is relevant.