1pub mod ai_udf;
5mod bridge;
6mod channel_source;
7pub mod complex_type_lambda;
9pub mod complex_type_udf;
11mod exec;
12pub mod execute;
14pub mod format_bridge_udf;
16pub mod json_extensions;
18pub mod json_path;
20pub mod json_tvf;
22pub mod json_types;
24pub mod json_udaf;
26pub mod json_udf;
28pub mod live_source;
30pub mod lookup_join;
32pub mod lookup_join_exec;
34pub mod proctime_udf;
36mod source;
37mod table_provider;
38pub mod watermark_udf;
41pub mod window_udf;
43
44pub use ai_udf::{ai_function_markers, AiFunctionMarker};
45pub use bridge::{BridgeSendError, BridgeSender, BridgeStream, BridgeTrySendError, StreamBridge};
46pub use channel_source::ChannelStreamSource;
47pub use complex_type_lambda::{
48 register_lambda_functions, ArrayFilter, ArrayReduce, ArrayTransform, MapFilter,
49 MapTransformValues,
50};
51pub use complex_type_udf::{
52 register_complex_type_functions, MapContainsKey, MapFromArrays, MapKeys, MapValues, StructDrop,
53 StructExtract, StructMerge, StructRename, StructSet,
54};
55pub use exec::StreamingScanExec;
56pub use execute::{execute_streaming_sql, DdlResult, QueryResult, StreamingSqlResult};
57pub use format_bridge_udf::{FromJsonUdf, ParseEpochUdf, ParseTimestampUdf, ToJsonUdf};
58pub use json_extensions::{
59 register_json_extensions, JsonInferSchema, JsonToColumns, JsonbDeepMerge, JsonbExcept,
60 JsonbFlatten, JsonbMerge, JsonbPick, JsonbRenameKeys, JsonbStripNulls, JsonbUnflatten,
61};
62pub use json_path::{CompiledJsonPath, JsonPathStep, JsonbPathExistsUdf, JsonbPathMatchUdf};
63pub use json_tvf::{
64 register_json_table_functions, JsonbArrayElementsTextTvf, JsonbArrayElementsTvf,
65 JsonbEachTextTvf, JsonbEachTvf, JsonbObjectKeysTvf,
66};
67pub use json_udaf::{JsonAgg, JsonObjectAgg};
68pub use json_udf::{
69 JsonBuildArray, JsonBuildObject, JsonTypeof, JsonbContainedBy, JsonbContains, JsonbExists,
70 JsonbExistsAll, JsonbExistsAny, JsonbGet, JsonbGetIdx, JsonbGetPath, JsonbGetPathText,
71 JsonbGetText, JsonbGetTextIdx, ToJsonb,
72};
73pub use live_source::{LiveSourceHandle, LiveSourceProvider};
74pub use lookup_join_exec::{
75 LookupJoinExec, LookupJoinExtensionPlanner, LookupSnapshot, LookupTableRegistry,
76 PartialLookupJoinExec, PartialLookupState, RegisteredLookup,
77};
78pub use proctime_udf::ProcTimeUdf;
79pub use source::{SortColumn, StreamSource, StreamSourceRef};
80pub use table_provider::StreamingTableProvider;
81pub use watermark_udf::WatermarkUdf;
82pub use window_udf::{
83 CumulateWindowEnd, CumulateWindowStart, HopWindowEnd, HopWindowStart, SessionWindowStart,
84 TumbleWindowEnd, TumbleWindowStart,
85};
86
87use std::sync::atomic::AtomicI64;
88use std::sync::Arc;
89
90use datafusion::execution::SessionStateBuilder;
91use datafusion::prelude::*;
92use datafusion_expr::ScalarUDF;
93
94use crate::planner::streaming_optimizer::{StreamingPhysicalValidator, StreamingValidatorMode};
95
96#[must_use]
103pub fn base_session_config() -> SessionConfig {
104 let mut config = SessionConfig::new();
105 config.options_mut().sql_parser.enable_ident_normalization = false;
106 config = config.with_target_partitions(1);
110 config
111}
112
113#[must_use]
122pub fn create_session_context() -> SessionContext {
123 SessionContext::new_with_config(base_session_config())
124}
125
126#[must_use]
150pub fn create_streaming_context() -> SessionContext {
151 create_streaming_context_with_validator(StreamingValidatorMode::Reject)
152}
153
154#[must_use]
162pub fn create_streaming_context_with_validator(mode: StreamingValidatorMode) -> SessionContext {
163 let config = base_session_config().with_batch_size(8192);
164
165 let ctx = if matches!(mode, StreamingValidatorMode::Off) {
166 SessionContext::new_with_config(config)
167 } else {
168 let default_state = SessionStateBuilder::new()
172 .with_config(config.clone())
173 .with_default_features()
174 .build();
175 let mut rules: Vec<
176 Arc<dyn datafusion::physical_optimizer::PhysicalOptimizerRule + Send + Sync>,
177 > = vec![Arc::new(StreamingPhysicalValidator::new(mode))];
178 rules.extend(default_state.physical_optimizers().iter().cloned());
179
180 let state = SessionStateBuilder::new()
181 .with_config(config)
182 .with_default_features()
183 .with_physical_optimizer_rules(rules)
184 .build();
185 SessionContext::new_with_state(state)
186 };
187
188 register_streaming_functions(&ctx);
189 ctx
190}
191
192fn register_non_watermark_udfs(ctx: &SessionContext) {
197 ctx.register_udf(ScalarUDF::new_from_impl(TumbleWindowStart::new()));
198 ctx.register_udf(ScalarUDF::new_from_impl(TumbleWindowEnd::new()));
199 ctx.register_udf(ScalarUDF::new_from_impl(HopWindowStart::new()));
200 ctx.register_udf(ScalarUDF::new_from_impl(HopWindowEnd::new()));
201 ctx.register_udf(ScalarUDF::new_from_impl(SessionWindowStart::new()));
202 ctx.register_udf(ScalarUDF::new_from_impl(CumulateWindowStart::new()));
203 ctx.register_udf(ScalarUDF::new_from_impl(CumulateWindowEnd::new()));
204 ctx.register_udf(ScalarUDF::new_from_impl(ProcTimeUdf::new()));
205 for marker in ai_function_markers() {
206 ctx.register_udf(marker);
207 }
208 register_json_functions(ctx);
209 register_json_extensions(ctx);
210 register_complex_type_functions(ctx);
211 register_lambda_functions(ctx);
212}
213
214pub fn register_streaming_functions(ctx: &SessionContext) {
219 register_non_watermark_udfs(ctx);
220 ctx.register_udf(ScalarUDF::new_from_impl(WatermarkUdf::unset()));
221}
222
223pub fn register_streaming_functions_with_watermark(
228 ctx: &SessionContext,
229 watermark_ms: Arc<AtomicI64>,
230) {
231 register_non_watermark_udfs(ctx);
232 ctx.register_udf(ScalarUDF::new_from_impl(WatermarkUdf::new(watermark_ms)));
233}
234
235pub fn register_json_functions(ctx: &SessionContext) {
238 ctx.register_udf(ScalarUDF::new_from_impl(JsonbGet::new()));
240 ctx.register_udf(ScalarUDF::new_from_impl(JsonbGetIdx::new()));
241 ctx.register_udf(ScalarUDF::new_from_impl(JsonbGetText::new()));
242 ctx.register_udf(ScalarUDF::new_from_impl(JsonbGetTextIdx::new()));
243 ctx.register_udf(ScalarUDF::new_from_impl(JsonbGetPath::new()));
244 ctx.register_udf(ScalarUDF::new_from_impl(JsonbGetPathText::new()));
245
246 ctx.register_udf(ScalarUDF::new_from_impl(JsonbExists::new()));
248 ctx.register_udf(ScalarUDF::new_from_impl(JsonbExistsAny::new()));
249 ctx.register_udf(ScalarUDF::new_from_impl(JsonbExistsAll::new()));
250
251 ctx.register_udf(ScalarUDF::new_from_impl(JsonbContains::new()));
253 ctx.register_udf(ScalarUDF::new_from_impl(JsonbContainedBy::new()));
254
255 ctx.register_udf(ScalarUDF::new_from_impl(JsonTypeof::new()));
257 ctx.register_udf(ScalarUDF::new_from_impl(JsonBuildObject::new()));
258 ctx.register_udf(ScalarUDF::new_from_impl(JsonBuildArray::new()));
259 ctx.register_udf(ScalarUDF::new_from_impl(ToJsonb::new()));
260
261 ctx.register_udaf(datafusion_expr::AggregateUDF::new_from_impl(JsonAgg::new()));
263 ctx.register_udaf(datafusion_expr::AggregateUDF::new_from_impl(
264 JsonObjectAgg::new(),
265 ));
266
267 ctx.register_udf(ScalarUDF::new_from_impl(ParseEpochUdf::new()));
269 ctx.register_udf(ScalarUDF::new_from_impl(ParseTimestampUdf::new()));
270 ctx.register_udf(ScalarUDF::new_from_impl(ToJsonUdf::new()));
271 ctx.register_udf(ScalarUDF::new_from_impl(FromJsonUdf::new()));
272
273 ctx.register_udf(ScalarUDF::new_from_impl(JsonbPathExistsUdf::new()));
275 ctx.register_udf(ScalarUDF::new_from_impl(JsonbPathMatchUdf::new()));
276
277 register_json_table_functions(ctx);
279}
280
281#[cfg(test)]
282mod tests;