A DataFrame's result as Arrow IPC streams, tech.ml.dataset datasets and dtype-next tensors, all at once or a batch at a time, and two previews: glimpse and to-html. Arrow, tech.ml.dataset and dtype-next are resolved when they're needed, so this loads with only Spark's Connect client, and without tech.ml.dataset.
A DataFrame's result as Arrow IPC streams, tech.ml.dataset datasets and dtype-next tensors, all at once or a batch at a time, and two previews: glimpse and to-html. Arrow, tech.ml.dataset and dtype-next are resolved when they're needed, so this loads with only Spark's Connect client, and without tech.ml.dataset.
(glimpse dataframe)(glimpse dataframe
{:keys [num-rows width] count? :count :or {num-rows 10 width 80}})Prints a transposed preview of dataframe: one line per column, with its
name, its type and its first few values, which reads better than show
on a wide DataFrame.
(g/glimpse df)
; Rows: at least 10
; Columns: 3
; $ id <bigint> 0, 1, 2, 3, 4, 5, 6, 7, 8, 9
; $ name <string> "Ada", "Bo", nil, "Grace", nil, "Alan", …
; $ score <double> 1.5, 2.5, 3.5, 4.5, 5.5, 6.5, 7.5, 8.5, …
Strings, numbers and nulls print as Clojure data, so a string is quoted and
a null is nil, and dates, times and other objects as their strings.
Rows: is exact only when the sample comes back short; otherwise it says
"at least", since counting the rows takes a full pass over the data.
Options:
:num-rows, the values to show per column, 10 by default;:width, where to cut a line, 80 by default, or ##Inf to cut nothing;:count, true to count the rows for an exact Rows:.Prints a transposed preview of `dataframe`: one line per column, with its name, its type and its first few values, which reads better than `show` on a wide DataFrame. ```clojure (g/glimpse df) ; Rows: at least 10 ; Columns: 3 ; $ id <bigint> 0, 1, 2, 3, 4, 5, 6, 7, 8, 9 ; $ name <string> "Ada", "Bo", nil, "Grace", nil, "Alan", … ; $ score <double> 1.5, 2.5, 3.5, 4.5, 5.5, 6.5, 7.5, 8.5, … ``` Strings, numbers and nulls print as Clojure data, so a string is quoted and a null is nil, and dates, times and other objects as their strings. `Rows:` is exact only when the sample comes back short; otherwise it says "at least", since counting the rows takes a full pass over the data. Options: - `:num-rows`, the values to show per column, 10 by default; - `:width`, where to cut a line, 80 by default, or `##Inf` to cut nothing; - `:count`, true to count the rows for an exact `Rows:`.
(stream dataframe)(stream dataframe {:keys [key-fn] :or {key-fn keyword}})The result of dataframe as tech.ml.dataset datasets, one per Arrow batch
that has rows, each as to-tmd makes it, with its options.
Returns a reducible, so that reduce, transduce, into and run! read
the batches as they go, and stop reading when they're done, stop early or
throw. It's also seqable, for seq, first and doseq, which read a
batch at a time but can stop before the end, so close it with with-open
for those:
(transduce (map tech.v3.dataset/row-count) + (g/stream df))
(with-open [batches (g/stream df)]
(first batches))
On classic Spark, each partition runs as a job of its own, as a reduce
gets to it, so only one partition's batches are on the driver at a time,
and a reduce is one of Spark's SQL executions, as collect is, so an
Observation from observe gets its metrics at its end, which a seq
doesn't give it. Over Spark Connect, the server sends the batches as it
makes them, and stopping releases the execution.
The result of `dataframe` as tech.ml.dataset datasets, one per Arrow batch that has rows, each as `to-tmd` makes it, with its options. Returns a reducible, so that `reduce`, `transduce`, `into` and `run!` read the batches as they go, and stop reading when they're done, stop early or throw. It's also seqable, for `seq`, `first` and `doseq`, which read a batch at a time but can stop before the end, so close it with `with-open` for those: ```clojure (transduce (map tech.v3.dataset/row-count) + (g/stream df)) (with-open [batches (g/stream df)] (first batches)) ``` On classic Spark, each partition runs as a job of its own, as a reduce gets to it, so only one partition's batches are on the driver at a time, and a reduce is one of Spark's SQL executions, as `collect` is, so an Observation from `observe` gets its metrics at its end, which a seq doesn't give it. Over Spark Connect, the server sends the batches as it makes them, and stopping releases the execution.
(stream-tensors dataframe)(stream-tensors dataframe {:keys [columns key-fn] :or {key-fn keyword}})The result of dataframe as maps of dtype-next tensors, one map per Arrow
batch that has rows, each as to-tensors makes it, with its options. Each
batch's tensors are its own: stacking them is up to the caller.
Returns a reducible, as stream does, which stops reading when a reduce is
done, stops early or throws. It's seqable too, a batch at a time, so close
it with with-open when reading it as a seq.
The result of `dataframe` as maps of dtype-next tensors, one map per Arrow batch that has rows, each as `to-tensors` makes it, with its options. Each batch's tensors are its own: stacking them is up to the caller. Returns a reducible, as `stream` does, which stops reading when a reduce is done, stops early or throws. It's seqable too, a batch at a time, so close it with `with-open` when reading it as a seq.
(to-arrow dataframe)The result of dataframe as Apache Arrow IPC streams in memory, one per
batch: a vector of byte arrays, each a complete stream of the schema, one
record batch and the end marker, for any Arrow reader. They must never be
concatenated byte for byte. An empty result gives one stream with no rows.
On classic Spark, the batches are Spark's own, as PySpark's toPandas
gets them, of at most spark.sql.execution.arrow.maxRecordsPerBatch rows
(10,000 by default). Over Spark Connect, they're the ones the server sends.
Like collect, the whole result comes to the driver.
The result of `dataframe` as Apache Arrow IPC streams in memory, one per batch: a vector of byte arrays, each a complete stream of the schema, one record batch and the end marker, for any Arrow reader. They must never be concatenated byte for byte. An empty result gives one stream with no rows. On classic Spark, the batches are Spark's own, as PySpark's `toPandas` gets them, of at most `spark.sql.execution.arrow.maxRecordsPerBatch` rows (10,000 by default). Over Spark Connect, they're the ones the server sends. Like `collect`, the whole result comes to the driver.
(to-html dataframe)(to-html dataframe {:keys [num-rows truncate] :or {num-rows 20 truncate 20}})The first rows of dataframe as an HTML table, as Spark renders a
DataFrame in a notebook: PySpark's _repr_html_. Spark escapes the cells,
and a note under the table says when there are more rows than it shows.
Options:
:num-rows, the rows to show, 20 by default;:truncate, the width that a cell is cut to, 20 by default, or 0 or
false not to cut.The first rows of `dataframe` as an HTML table, as Spark renders a DataFrame in a notebook: PySpark's `_repr_html_`. Spark escapes the cells, and a note under the table says when there are more rows than it shows. Options: - `:num-rows`, the rows to show, 20 by default; - `:truncate`, the width that a cell is cut to, 20 by default, or 0 or false not to cut.
(to-tensors dataframe)(to-tensors dataframe {:keys [columns key-fn] :or {key-fn keyword}})The result of dataframe as dtype-next tensors, one per column, in a map
by column name: by keyword, or by :key-fn applied to the name.
:columns selects the columns first. The whole result comes to the
driver; stream-tensors is for one batch at a time.
A column of integers (TINYINT, SMALLINT, INT or BIGINT) or floating-point
numbers (FLOAT or DOUBLE) with no nulls becomes a tensor of shape [rows],
and a column of arrays of those, or of dense MLlib vectors, all of one
length, a tensor of shape [rows length]. Anything else throws, naming the
column: nulls, DECIMALs, strings, booleans, arrays of different lengths
and sparse vectors. So do an empty result, rows without columns, and two
columns that :key-fn names alike.
(g/to-tensors scored {:columns [:features :label]})
;; => {:features #tech.v3.tensor<float64>[1000 4] ..., :label ...}
It needs dtype-next on the classpath, which tech.ml.dataset brings, and
over Spark Connect, org.apache.arrow/arrow-vector and
arrow-memory-netty.
The result of `dataframe` as dtype-next tensors, one per column, in a map
by column name: by keyword, or by `:key-fn` applied to the name.
`:columns` selects the columns first. The whole result comes to the
driver; `stream-tensors` is for one batch at a time.
A column of integers (TINYINT, SMALLINT, INT or BIGINT) or floating-point
numbers (FLOAT or DOUBLE) with no nulls becomes a tensor of shape [rows],
and a column of arrays of those, or of dense MLlib vectors, all of one
length, a tensor of shape [rows length]. Anything else throws, naming the
column: nulls, DECIMALs, strings, booleans, arrays of different lengths
and sparse vectors. So do an empty result, rows without columns, and two
columns that `:key-fn` names alike.
```clojure
(g/to-tensors scored {:columns [:features :label]})
;; => {:features #tech.v3.tensor<float64>[1000 4] ..., :label ...}
```
It needs dtype-next on the classpath, which tech.ml.dataset brings, and
over Spark Connect, `org.apache.arrow/arrow-vector` and
`arrow-memory-netty`.(to-tmd dataframe)(to-tmd dataframe {:keys [key-fn] :or {key-fn keyword}})The result of dataframe as one tech.ml.dataset dataset, which needs
techascent/tech.ml.dataset on the classpath. Like collect, the whole
result comes to the driver; stream is for one batch at a time. An empty
result gives a dataset with no rows and the result's columns.
Columns are named by keyword, or by :key-fn applied to the name. A null
is a missing value. Numbers and booleans stay primitive. DECIMAL becomes a
BigDecimal, DATE a LocalDate, TIMESTAMP an Instant, TIMESTAMP_NTZ a
LocalDateTime, TIME a LocalTime, a day-time interval a Duration, a
year-month interval a Period and BINARY a byte array. An array becomes a
vector, and a struct or a map a map, with keyword keys for a struct's
fields. VARIANT becomes Spark's VariantVal, and an MLlib vector what
collect gives. Each column keeps its Spark type, as DDL, in its metadata
under :zero-one.geni/spark-type, which create-dataframe uses, but for
one that holds MLlib vectors.
A calendar interval, a geometry or a geography throws, as do a struct with
two fields of one name, two columns of one name, and two that :key-fn
names alike, before a job runs. So do rows without columns, which a
dataset can't hold. On classic Spark, it runs as one of Spark's SQL
executions, as collect does, so an Observation from observe gets its
metrics.
Over Spark Connect, it needs org.apache.arrow/arrow-vector and
arrow-memory-netty on the classpath, since the client's Arrow is shaded.
The result of `dataframe` as one tech.ml.dataset dataset, which needs `techascent/tech.ml.dataset` on the classpath. Like `collect`, the whole result comes to the driver; `stream` is for one batch at a time. An empty result gives a dataset with no rows and the result's columns. Columns are named by keyword, or by `:key-fn` applied to the name. A null is a missing value. Numbers and booleans stay primitive. DECIMAL becomes a BigDecimal, DATE a LocalDate, TIMESTAMP an Instant, TIMESTAMP_NTZ a LocalDateTime, TIME a LocalTime, a day-time interval a Duration, a year-month interval a Period and BINARY a byte array. An array becomes a vector, and a struct or a map a map, with keyword keys for a struct's fields. VARIANT becomes Spark's VariantVal, and an MLlib vector what `collect` gives. Each column keeps its Spark type, as DDL, in its metadata under `:zero-one.geni/spark-type`, which `create-dataframe` uses, but for one that holds MLlib vectors. A calendar interval, a geometry or a geography throws, as do a struct with two fields of one name, two columns of one name, and two that `:key-fn` names alike, before a job runs. So do rows without columns, which a dataset can't hold. On classic Spark, it runs as one of Spark's SQL executions, as `collect` does, so an Observation from `observe` gets its metrics. Over Spark Connect, it needs `org.apache.arrow/arrow-vector` and `arrow-memory-netty` on the classpath, since the client's Arrow is shaded.
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 |