Spark SQL UDFs from Clojure functions, on classic Spark.
Spark SQL UDFs from Clojure functions, on classic Spark.
Registers f as a Spark UDF called udf-name on the session, for
g/sql and g/expr, and returns a function of columns that calls it, as
g/udf does. The return type and the options are as for g/udf, plus
:arity: the number of columns it takes, for a function that takes more
than one number of arguments, or any number of them.
(g/register-udf! "plus_one" inc :long)
(-> (g/range 3)
(g/select {:x (g/expr "plus_one(id)")})
g/collect)
=> ({:x 1} {:x 2} {:x 3})
Registers `f` as a Spark UDF called `udf-name` on the session, for
`g/sql` and `g/expr`, and returns a function of columns that calls it, as
`g/udf` does. The return type and the options are as for `g/udf`, plus
`:arity`: the number of columns it takes, for a function that takes more
than one number of arguments, or any number of them.
```clojure
(g/register-udf! "plus_one" inc :long)
(-> (g/range 3)
(g/select {:x (g/expr "plus_one(id)")})
g/collect)
=> ({:x 1} {:x 2} {:x 3})
```(udf f return-type)(udf f return-type opts)Returns a function of columns that applies f to their values, one row at
a time, as a Spark UDF, and returns the result as a Column.
f gets each row's values as Clojure data, as g/collect gives them: nil
for a null, a seq for an array, a map for a map or a struct. Its result is
converted to return-type, which is a type keyword such as :long or :string,
a schema in g/->schema's form, such as [:string] for an array of strings
or {:a :int} for a struct, or a Spark DataType. So a long becomes an int for
:int, and a map becomes a struct for {:a :int}.
The options are:
:name, which the column's name and g/explain show;:deterministic, false when f can return different results for the same
values, so that Spark calls it once per row, as it's written;:nullable, false when f never returns nil.UDFs need classic Spark, and run on the executors. Functions defined at a
REPL, or in a script, work on a local session that Geni starts. On a
cluster, pass a var, such as #'my-fn, which the executors look up in its
namespace, or AOT-compile the namespace that defines f. The Clojure UDFs
guide has the details.
(def plus-one (g/udf inc :long))
(-> (g/range 3)
(g/select {:x (plus-one :id)})
g/collect)
=> ({:x 1} {:x 2} {:x 3})
Returns a function of columns that applies `f` to their values, one row at
a time, as a Spark UDF, and returns the result as a Column.
`f` gets each row's values as Clojure data, as `g/collect` gives them: nil
for a null, a seq for an array, a map for a map or a struct. Its result is
converted to `return-type`, which is a type keyword such as :long or :string,
a schema in `g/->schema`'s form, such as [:string] for an array of strings
or {:a :int} for a struct, or a Spark DataType. So a long becomes an int for
:int, and a map becomes a struct for {:a :int}.
The options are:
- `:name`, which the column's name and `g/explain` show;
- `:deterministic`, false when `f` can return different results for the same
values, so that Spark calls it once per row, as it's written;
- `:nullable`, false when `f` never returns nil.
UDFs need classic Spark, and run on the executors. Functions defined at a
REPL, or in a script, work on a local session that Geni starts. On a
cluster, pass a var, such as `#'my-fn`, which the executors look up in its
namespace, or AOT-compile the namespace that defines `f`. The Clojure UDFs
guide has the details.
```clojure
(def plus-one (g/udf inc :long))
(-> (g/range 3)
(g/select {:x (plus-one :id)})
g/collect)
=> ({:x 1} {:x 2} {:x 3})
```cljdoc builds & hosts documentation for Clojure/Script libraries
| Ctrl+k | Jump to recent docs |
| ← | Move to previous article |
| → | Move to next article |
| Ctrl+/ | Jump to the search field |