Spark SQL UDFs from Clojure functions, on classic Spark and over Spark Connect.
Spark SQL UDFs from Clojure functions, on classic Spark and over Spark Connect.
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, as one that draws random numbers can, so that Spark's optimiser
doesn't move it or merge it with other expressions as it can a
deterministic one. It says nothing about how many times Spark calls f
for a row: a retried task, or a DataFrame that's computed twice, calls it
again, so any side effects have to cope with repeated calls;:nullable, false when f never returns nil.UDFs 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. Over Spark Connect, Geni uploads Clojure, Geni
and the code that f uses to the server, once per session, and the server
loads each namespace that has a file from that file. A function defined at
the REPL works when it's defined after (g/connect url {:keep-classes true}). Code that changed after it went to a session throws an error, since
the server keeps what it got first. 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, as one that draws random numbers can, so that Spark's optimiser
doesn't move it or merge it with other expressions as it can a
deterministic one. It says nothing about how many times Spark calls `f`
for a row: a retried task, or a DataFrame that's computed twice, calls it
again, so any side effects have to cope with repeated calls;
- `:nullable`, false when `f` never returns nil.
UDFs 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`. Over Spark Connect, Geni uploads Clojure, Geni
and the code that `f` uses to the server, once per session, and the server
loads each namespace that has a file from that file. A function defined at
the REPL works when it's defined after `(g/connect url {:keep-classes
true})`. Code that changed after it went to a session throws an error, since
the server keeps what it got first. 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 |