risingwave_expr_impl/scalar/
date_bin.rs1use risingwave_common::types::{Interval, Timestamp, Timestamptz};
16use risingwave_expr::{ExprError, Result, function};
17
18#[function("date_bin(interval, timestamp, timestamp) -> timestamp")]
19pub fn date_bin_ts(stride: Interval, source: Timestamp, origin: Timestamp) -> Result<Timestamp> {
20 let source_us = source.0.and_utc().timestamp_micros(); let origin_us = origin.0.and_utc().timestamp_micros(); let binned_source_us = date_bin_inner(stride, source_us, origin_us)?;
24 Timestamp::with_micros(binned_source_us).map_err(|_| ExprError::NumericOutOfRange)
25}
26
27#[function("date_bin(interval, timestamptz, timestamptz) -> timestamptz")]
28pub fn date_bin_tstz(
29 stride: Interval,
30 source: Timestamptz,
31 origin: Timestamptz,
32) -> Result<Timestamptz> {
33 let source_us = source.timestamp_micros(); let origin_us = origin.timestamp_micros(); let binned_source_us = date_bin_inner(stride, source_us, origin_us)?;
37 Timestamptz::from_micros(binned_source_us).ok_or(ExprError::NumericOutOfRange)
38}
39
40fn date_bin_inner(stride: Interval, source_us: i64, origin_us: i64) -> Result<i64> {
41 if stride.months() != 0 {
42 return Err(ExprError::InvalidParam {
44 name: "stride",
45 reason: "stride interval with months not supported in date_bin".into(),
46 });
47 }
48 let stride_us = stride.usecs() + (stride.days() as i64) * Interval::USECS_PER_DAY; if stride_us <= 0 {
51 return Err(ExprError::InvalidParam {
52 name: "stride",
53 reason: "stride interval must be positive".into(),
54 });
55 }
56
57 let delta = source_us - origin_us;
59
60 let bucket = delta.div_euclid(stride_us) * stride_us;
62
63 let binned_source_us = origin_us + bucket;
65 Ok(binned_source_us)
66}