When using ElasticSearch, how can I bucket event-driven data over time and deal with different update cadences?

Viewed 51

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:

  1. It isn't just 1 balance field, I have dozens of these fields.
  2. 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
  3. 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.

0 Answers
Related