Spark functions run on the cluster, but sometimes it's easier to work on the data in plain Clojure. That means moving the data from the Spark workers to the driver, where the Clojure REPL runs, which only works when the data is small enough to fit on the driver.
Geni's functions that start with collect bring the data to the driver. The examples below use the Melbourne housing data in Geni's repo:
(require '[zero-one.geni.core :as g])
(def dataframe
(-> (g/read-parquet! "test/resources/melbourne_housing_snapshot.parquet")
(g/select :Suburb :Address :Rooms :Price :Date)))
A very common case is to access the data as Clojure maps with collect:
(-> dataframe (g/limit 2) g/collect)
;; => ({:Suburb "Abbotsford",
;; :Address "85 Turner St",
;; :Rooms 2,
;; :Price 1480000.0,
;; :Date "3/12/2016"}
;; {:Suburb "Abbotsford",
;; :Address "25 Bloomburg St",
;; :Rooms 2,
;; :Price 1035000.0,
;; :Date "4/02/2016"})
Alternatively, collect-vals returns a sequence of vectors:
(-> dataframe (g/limit 2) g/collect-vals)
;; => (["Abbotsford" "85 Turner St" 2 1480000.0 "3/12/2016"]
;; ["Abbotsford" "25 Bloomburg St" 2 1035000.0 "4/02/2016"])
To access the values of a single column, use collect-col:
(-> dataframe (g/limit 2) (g/collect-col :Address))
;; => ("85 Turner St" "25 Bloomburg St")
collect-to-arrow brings the data to the driver as Arrow files instead. It can handle data larger than the driver's heap, as long as the largest partition fits on the driver, since the data travels one partition at a time. Repartitioning the data first makes sure of that.
It also needs to know how many rows each Arrow file gets, which should be small enough for each file to fit in the heap, and the directory to write the files to. It writes files of chunk-size rows each (the last one can be smaller) and returns their paths:
(-> dataframe
(g/repartition 20) ;; 20 partitions of about the same size
(g/collect-to-arrow 1000 "/tmp"))
; ["/tmp/geni12331590604347994819.ipc"
; "/tmp/geni2107925719499812901.ipc"
; ...]
With enough partitions and a small enough chunk size, data of any size can make it to the driver, although slowly when there's a lot of it.
The files are in the Arrow streaming format, which other tools can read, such as tech.ml.dataset with tech.v3.libs.arrow/stream->dataset. On Spark 3.5 with JDK 21 or newer, collect-to-arrow needs Arrow 13 or newer on the classpath (see the installation notes).
tech.ml.dataset also has a tech.v3.libs.spark namespace, which converts a Spark dataset into a tech.ml.dataset dataset on the driver, and back. The data has to fit in the driver's heap.
Can you improve this documentation? These fine people already did:
Anthony Khong & Carsten BehringEdit on GitHub
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 |