Repository navigation
Feedback request for providing configurable UDF functions #10744
Description
Activity
- changed the title
[-]Switch UDF function expression definitions from static instance calls to use function_registry for lookup[/-][+]Feedback request for providing configurable UDF functions[/+]on May 31, 2024 I think options 1 and 3 would be straightforward
You could even potentially implement
pub fn to_timestamp_safe(args: Vec<Expr>) -> Expr { ... }
Directly in your application (rather than in the core of datafusion)
Another crazy thought might be to implement a rewrite pass (e.g.
AnalyzerRule) that rewrites all expressions in the plan when safe mode is needed... I think they have access to all the state necessaryI think the key thing to figure out is "will safemode to_timestamp be part of the datafusion core"?
Maybe it is time to make a
datafusion-functions-scalar-sparkor something that has functions that have the spark behaviors 🤔I think it is possible to extend the
safe_modeto the datafusion core like what you mentioned, it should be similar to thedistinctmode to the aggregate function. We can have different behavior based on whethersafe_modeis set.For the expression API, we can either
- Introduce
to_timestamp_safe() - Extend
to_timestamp(args, safe: bool) - Introduce builder mode to avoid breaking change
to_timestamp(args).safe().build()
I prefer the third one.
Also, there are many
to_timestampfunctions with different time units. I'm not sure why they are split into different functions, they are possible to collapse into a single to_timestamp function and cast to different time units based on the given argument.We can have
to_timestamp(args).time_unit(Mili).safe().build()if we need the timestamp millisecond andto_timestamp(args).safe().build()for Nanosecond"will safemode to_timestamp be part of the datafusion core"
It is an interesting question, we can think of implementing functions based on other DB in the first place.
For example, we usually follow postgres, duckdb, and others.
We can have
functions-postgres,functions-duckdbandfunctions-sparkthat aim to mirror the behavior of other db, and we don't even needfunctionscrate (we can keep them for backward compatibility). They are considered as theextension crate implemented by the third party, datafusion does not need to implement any datafusion-specific function (and we prefer not to), we just make sure datafusion core is possible compatible with different functions crate. And we can register most-used functions to the datafusion for the end-to-end SQL workflow!- Introduce
We can have functions-postgres, functions-duckdb and functions-spark
Most of the function has the same behavior in different db, we can also implement different functions in one crate
functions, but we register the expected functions based on the DB we want to mimic. Instead of having adefault functionfor datafusion SQL workflow, we can havespark_functions, andpostgres_functionsthat register differentto_timestampfunctions based on the configuration we set.We can have functions-postgres, functions-duckdb and functions-spark
Most of the function has the same behavior in different db, we can also implement different functions in one crate
functions, but we register the expected functions based on the DB we want to mimic. Instead of having adefault functionfor datafusion SQL workflow, we can havespark_functions, andpostgres_functionsthat register differentto_timestampfunctions based on the configuration we set.@andygrove are there udf's already in the comet project that handle spark specific behaviour? If so is that a separate project or embedded in comet currently? (I haven't looked at that codebase myself since the initial drop)
The one issue with moving this functionality into a spark module is that for that to really be valid the formats would have to be spark compatible, which they are not currently. I do not have the spare time in the near future to implement a parser to do that.
Reacted by Andrew Lamb@Omega359 so far we have been implementing custom
PhysicalExprdirectly in thedatafusion-cometproject as needed for Spark-specific behavior, with support for Spark's different evaluation modes (legacy, try, ansi) and we are using fuzz testing to ensure compatibilty across multiple Spark versions.I think we need to have the discussion of whether it makes sense to upstream these into the core datafusion project or not, or whether we publish a
datafusion-spark-compatcrate from Comet, or some other option.The one issue with moving this functionality into a spark module is that for that to really be valid the formats would have to be spark compatible, which they are not currently. I do not have the spare time in the near future to implement a parser to do that.
We are porting Spark parsing logic as part of Comet.
I think we need to have the discussion of whether it makes sense to upstream these into the core datafusion project or not, or whether we publish a
datafusion-spark-compatcrate from Comet, or some other option.Thank you for chiming in. While I wouldn't mind spark compatibility it really isn't the focus of this request as I've already converted all the spark expressions and function usages to DF compatible ones. It's the general system behaviour that is what I would like to address - being able to essentially switch from a db focused perspective (fail fast) to a processing engine one (nominally lenient - return null) for some (all) of the UDF's.
If the general consensus is to separate out this desired behaviour than I would think a separate crate might be the best approach. However from searching the issues here there seems to have been some talk of how to handle mirroring the behaviour of other databases in the past but it also includes sql syntax as well so it's not quite as simple as just having a db specific crate full of UDF's and calling it a day.
I have read the context now and understand that this is about
safemode or what Spark callsANSImode.Isn't this just a case of adding a new flag to the session context that UDFs can choose to use when deciding whether to return null or throw an error?
I have read the context now and understand that this is about
safemode or what Spark callsANSImode.Isn't this just a case of adding a new flag to the session context that UDFs can choose to use when deciding whether to return null or throw an error?
That would be nice ... except UDF's don't have a way to access the session context currently :( Option #2 and #3 provide that via different mechanisms.
I wonder if we could take a page from what @jayzhan211 is implementing in #10560 and go with a trait
So we could implement something like
let expr = to_timestamp(lit("2021-01-01")) // set the to_timestamp mode to "safe" .safe();
I realize that this would require changing the callsites so maybe it isn't viable
After thinking about this a fair bit the builder approach like what @jayzhan211 did with aggregate functions seems to be the best way forward on this feature imho. While I do like the idea of a separate crate(s) for mirroring functionality from other systems I think that is a much much larger project and is encompasses a lot more functionality than this specific feature entails. Putting this feature into core I don't believe limits DF in the future to extracting out this and other similar behaviour 'traits' and functionality to system specific crates.
I'll start work on this and see how that works out. If it does then I'll add safe support via a trait to the to_timestamp*, to_date and to_unixtime functions. If there are other UDF's that could benefit from having a 'safe' mode (return null on error) please let me know and I'll see about adding safe mode to those as well.
Thank you everyone for your feedback and guidance on this feedback request! 👍
Reacted by Jay ZhanWe just merged the aggregate builder in #10560 -- I am quite happy with how it turned out, in case you want to take a friendly look
Reacted by Bruce RitchieAfter attempting to implement the builder approach it became apparent to me that it will touch too many things and really won't work well without changing the signature of ScalarUDFImpl anyways. It works for the aggregate functions because the functions defined in the AggregateUDFImpl trait have arguments where the additional information (distinct, sort, ordering, etc) is provided to the UDF implementation. In the case of ScalarUDFImpl though that is not the case.
After some more thought I think the cleanest approach may be to add a get_config_options function to the SimplifyInfo trait and add a
scalar_udf_safe_mode: bool, default = falseto the ExecutionOptions struct. Doing that will allow functions that require configuration (including but obviously not limited to the 'safe' mode I'm working on) to access them while changing as little as possible wrt trait signatures.err, scratch that. Onto the next idea :/
Reacted by Andrew Lamb
Is your feature request related to a problem or challenge?
During work on adding a 'safe' mode to to_timestamp and to_date UDF functions I've come across an issue that I would like feedback on before proceeding.
The feature
Currently for timestamp and date parsing if a source string cannot be parsed using any of the provided chrono formats datafusion will return an error. This is normal for a database-type solution however it is not ideal for a system that is parsing billions of dates in batches - some of which are human entered. Systems such as Spark default to a null value for anything that cannot be parsed and this feature enables a mode ('safe' to mirror the name and same behaviour as CastOptions) that allows the to_timestamp* and to_date UDFs to have the same behaviour (return null on error).
The problem
Since UDF functions have no context provided and there isn't a way I know of to statically get access to config to add the above mentioned functionality I resorted to using a new constructor function to allow the UDF to switch behaviour:
To use the alternative 'safe' mode for these functions is as simple as
Unfortunately this only affect sql queries - any calls to the to_timestamp(args: Vec) function will not use the new definition as registered in the function registry. This is because that function and every other function like it use a static singleton instance that only uses a ::new() call to initial it and there is no way that I can see to replace that instance.
Describe the solution you'd like
I see a few possible solutions to this:
ctx.udf("to_timestamp").unwrap().call(args)instead of the to_timestamp() function anytime 'safe' mode is required. This is less than ideal imho as it can lead to confusion and unintuitive behavior.Any opinions, suggestions and critiques would be welcome.
Describe alternatives you've considered
No response
Additional context
No response