feat(dataframe): implement missing methods (#67)

* feat(dataframe): stat & na methods

* update README

* update rust docs

* flaky test
4 files changed
tree: d1b748d985ce23be0455eef0f19ffc12d95bfb9b
  1. .github/
  2. core/
  3. datasets/
  4. examples/
  5. rust/
  6. .gitignore
  7. .pre-commit-config.yaml
  8. Cargo.lock
  9. Cargo.toml
  10. docker-compose.yml
  11. LICENSE.txt
  12. pre-commit.sh
  13. README.md
README.md

Apache Spark Connect Client for Rust

This project houses the experimental client for Spark Connect for Apache Spark written in Rust

Current State of the Project

Currently, the Spark Connect client for Rust is highly experimental and should not be used in any production setting. This is currently a “proof of concept” to identify the methods of interacting with Spark cluster from rust.

The spark-connect-rs aims to provide an entrypoint to Spark Connect, and provide similar DataFrame API interactions.

Project Layout

├── core         <- core implementation in Rust
│   └─ protobuf  <- connect protobuf for apache/spark
├── rust         <- shim for 'spark-connect-rs' from core
├── examples     <- examples of using different aspects of the crate
├── datasets     <- sample files from the main spark repo

Future state would be to have additional bindings for other languages along side the top level rust folder.

Getting Started

This section explains how run Spark Connect Rust locally starting from 0.

Step 1: Install rust via rustup: https://www.rust-lang.org/tools/install

Step 2: Ensure you have a cmake and protobuf installed on your machine

Step 3: Run the following commands to clone the repo

git clone https://github.com/sjrusso8/spark-connect-rs.git

cargo build

Step 4: Setup the Spark Driver on localhost either by downloading spark or with docker.

With local spark:

  1. Download Spark distribution (3.5.1 recommended), unzip the package.

  2. Set your SPARK_HOME environment variable to the location where spark was extracted to,

  3. Start the Spark Connect server with the following command (make sure to use a package version that matches your Spark distribution):

$ $SPARK_HOME/sbin/start-connect-server.sh --packages "org.apache.spark:spark-connect_2.12:3.5.1,io.delta:delta-spark_2.12:3.0.0" \
      --conf "spark.driver.extraJavaOptions=-Divy.cache.dir=/tmp -Divy.home=/tmp" \
      --conf "spark.sql.extensions=io.delta.sql.DeltaSparkSessionExtension" \
      --conf "spark.sql.catalog.spark_catalog=org.apache.spark.sql.delta.catalog.DeltaCatalog"

With docker:

  1. Start the Spark Connect server by leveraging the created docker-compose.yml in this repo. This will start a Spark Connect Server running on port 15002
$ docker compose up --build -d

Step 5: Run an example from the repo under /examples

Features

The following section outlines some of the larger functionality that are not yet working with this Spark Connect implementation.

  • done TLS authentication & Databricks compatability via the feature flag feature = 'tls'
  • open StreamingQueryManager
  • open UDFs or any type of functionality that takes a closure (foreach, foreachBatch, etc.)

SparkSession

Spark Session type object and its implemented traits

SparkSessionAPIComment
activeopen
addArtifact(s)open
addTagdone
clearTagsdone
copyFromLocalToFsopen
createDataFramepartialPartial. Only works for RecordBatch
getActiveSessionsopen
getTagsdone
interruptAlldonesplitn
interruptOperationdone
interruptTagdone
newSessionopen
rangedone
removeTagdone
sqldone
stopopen
tabledone
catalogdoneCatalog
clientdoneunstable developer api for testing only
confdoneConf
readdoneDataFrameReader
readStreamdoneDataStreamReader
streamsopenStreams
udfopenUdf - may not be possible
udtfopenUdtf - may not be possible
versiondone

SparkSessionBuilder

SparkSessionBuilderAPIComment
appNamedone
configdone
masteropen
remotepartialValidate using spark connection string

StreamingQueryManager

StreamingQueryManagerAPIComment
awaitAnyTerminationopen
getopen
resetTerminatedopen
activeopen

StreamingQuery

StreamingQueryAPIComment
awaitTerminationdone
exceptiondone
explaindone
processAllAvailabledone
stopdone
iddone
isActivedone
lastProgressdone
namedone
recentProgressdone
runIddone
statusdone

DataStreamReader

DataStreamReaderAPIComment
csvopen
formatdone
jsonopen
loaddone
optiondone
optionsdone
orcopen
parquetopen
schemadone
tableopen
textopen

DataFrameReader

DataFrameReaderAPIComment
csvopen
formatdone
jsonopen
loaddone
optiondone
optionsdone
orcopen
parquetopen
schemadone
tabledone
textopen

DataStreamWriter

Start a streaming job and return a StreamingQuery object to handle the stream operations.

DataStreamWriterAPIComment
foreach
foreachBatch
formatdone
optiondone
optionsdone
outputModedoneUses an Enum for OutputMode
partitionBydone
queryNamedone
startdone
toTabledone
triggerdoneUses an Enum for TriggerMode

StreamingQueryListener

StreamingQueryListenerAPIComment
onQueryIdleopen
onQueryProgressopen
onQueryStartedopen
onQueryTerminatedopen

UdfRegistration (may not be possible)

UDFRegistrationAPIComment
registeropen
registerJavaFunctionopen
registerJavaUDAFopen

UdtfRegistration (may not be possible)

UDTFRegistrationAPIComment
registeropen

RuntimeConfig

RuntimeConfigAPIComment
getdone
isModifiabledone
setdone
unsetdone

Catalog

CatalogAPIComment
cacheTabledone
clearCachedone
createExternalTaleopen
createTableopen
currentCatalogdone
currentDatabasedone
databaseExistsdone
dropGlobalTempViewdone
dropTempViewdone
functionExistsdone
getDatabasedone
getFunctiondone
getTabledone
isCacheddone
listCatalogsdone
listDatabasesdone
listFunctionsdone
listTablesdone
recoverPartitionsdone
refreshByPathdone
refreshTabledone
registerFunctionopen
setCurrentCatalogdone
setCurrentDatabasedone
tableExistsdone
uncacheTabledone

DataFrame

Spark DataFrame type object and its implemented traits.

DataFrameAPIComment
aggdone
aliasdone
approxQuantiledone
cachedone
checkpointopenNot part of Spark Connect
coalescedone
colRegexdone
collectdone
columnsdone
corrdone
countdone
covdone
createGlobalTempViewdone
createOrReplaceGlobalTempViewdone
createOrReplaceTempViewdone
createTempViewdone
crossJoindone
crosstabdone
cubedone
describedone
distinctdone
dropdone
dropDuplicatesdone
dropDuplicatesWithinWatermarkdone
drop_duplicatesdone
dropnadone
dtypesdone
exceptAlldone
explaindone
fillnadone
filterdone
firstdone
foreachopen
foreachPartitionopen
freqItemsdone
groupBydone
headdone
hintdone
inputFilesdone
intersectdone
intersectAlldone
isEmptydone
isLocaldone
isStreamingdone
joindone
limitdone
localCheckpointopenNot part of Spark Connect
mapInPandasopenTBD on this exact implementation
mapInArrowopenTBD on this exact implementation
meltdone
nadone
observeopen
offsetdone
orderBydone
persistdone
printSchemadone
randomSplitdone
registerTempTabledone
repartitiondone
repartitionByRangedone
replacedone
rollupdone
sameSemanticsdone
sampledone
sampleBydone
schemadone
selectdone
selectExprdone
semanticHashdone
showdone
sortdone
sortWithinPartitionsdone
sparkSessiondone
statdone
storageLeveldone
subtractdone
summarydone
taildone
takedone
todone
toDFdone
toJSONpartialDoes not return an RDD but a long JSON formatted String
toLocalIteratoropen
toPandas to_polars & toPolarspartialConvert to a polars::frame::DataFrame
new to_datafusion & toDataFusiondoneConvert to a datafusion::dataframe::DataFrame
transformdone
uniondone
unionAlldone
unionByNamedone
unpersistdone
unpivotdone
wheredoneuse filter instead, where is a keyword for rust
withColumndone
withColumnsdone
withColumnRenameddone
withColumnsRenameddone
withMetadatadone
withWatermarkdone
writedone
writeStreamdone
writeTodone

DataFrameWriter

Spark Connect should respect the format as long as your cluster supports the specified type and has the required jars

DataFrameWriterAPIComment
bucketBydone
csv
formatdone
insertIntodone
jdbc
json
modedone
optiondone
optionsdone
orc
parquet
partitionBy
savedone
saveAsTabledone
sortBydone
text

DataFrameWriterV2

DataFrameWriterV2APIComment
appenddone
createdone
createOrReplacedone
optiondone
optionsdone
overwritedone
overwritePartitionsdone
partitionedBydone
replacedone
tablePropertydone
usingdone

Column

Spark Column type object and its implemented traits

ColumnAPIComment
aliasdone
ascdone
asc_nulls_firstdone
asc_nulls_lastdone
astypeopen
betweenopen
castdone
containsdone
descdone
desc_nulls_firstdone
desc_nulls_lastdone
dropFieldsdone
endswithdone
eqNullSafeopen
getFieldopenThis is depreciated but will need to be implemented
getItemopenThis is depreciated but will need to be implemented
ilikedone
isNotNulldone
isNulldone
isindone
likedone
namedone
otherwiseopen
overdoneRefer to Window for creating window specifications
rlikedone
startswithdone
substrdone
whenopen
withFielddone
eq ==doneRust does not like when you try to overload == and return something other than a bool. Currently implemented column equality like col('name').eq(col('id')). Not the best, but it works for now
addition +done
subtration -done
multiplication *done
division /done
OR |done
AND &done
XOR ^done
Negate ~done

Functions

Only a few of the functions are covered by unit tests.

FunctionsAPIComment
absdone
acosdone
acoshdone
add_monthsdone
aggregateopen
approxCountDistinctopen
approx_count_distinctdone
arraydone
array_appenddone
array_compactdone
array_containsopen
array_distinctdone
array_exceptdone
array_insertopen
array_intersectdone
array_joinopen
array_maxdone
array_mindone
array_positiondone
array_removedone
array_repeatdone
array_sortopen
array_uniondone
arrays_overlapopen
arrays_zipdone
ascdone
asc_nulls_firstdone
asc_nulls_lastdone
asciidone
asindone
asinhdone
assert_trueopen
atandone
atan2done
atanhdone
avgdone
base64done
bindone
bit_lengthdone
bitwiseNOTopen
bitwise_notdone
broadcastopen
broundopen
bucketopen
call_udfopen
cbrtdone
ceildone
coalescedone
coldone
collect_listdone
collect_setdone
columndone
concatdone
concat_wsopen
convopen
corropen
cosopen
coshopen
cotopen
countopen
countDistinctopen
count_distinctopen
covar_popdone
covar_sampdone
crc32done
create_mapdone
cscdone
cume_distdone
current_datedone
current_timestampdone
date_adddone
date_formatopen
date_subdone
date_truncopen
datediffdone
dayofmonthdone
dayofweekdone
dayofyeardone
daysdone
decodeopen
degreesdone
dense_rankdone
descdone
desc_nulls_firstdone
desc_nulls_lastdone
element_atopen
encodeopen
existsopen
expdone
explodedone
explode_outerdone
expm1done
exprdone
factorialdone
filteropen
firstopen
flattendone
floordone
forallopen
format_numberopen
format_stringopen
from_csvopen
from_jsonopen
from_unixtimeopen
from_utc_timestampopen
functoolsopen
getopen
get_active_spark_contextopen
get_json_objectopen
greatestdone
groupingdone
grouping_idopen
has_numpyopen
hashdone
hexdone
hourdone
hoursdone
hypotopen
initcapdone
inlinedone
inline_outerdone
input_file_namedone
inspectopen
instropen
isnandone
isnulldone
json_tupleopen
kurtosisdone
lagopen
lastopen
last_dayopen
leadopen
leastdone
lengthdone
levenshteinopen
litdone
localtimestampdone
locateopen
logdone
log10done
log1pdone
log2done
lowerdone
lpadopen
ltrimdone
make_dateopen
map_concatdone
map_contains_keyopen
map_entriesdone
map_filteropen
map_from_arraysopen
map_from_entriesdone
map_keysdone
map_valuesdone
map_zip_withopen
maxdone
max_byopen
md5done
meandone
mediandone
mindone
min_byopen
minutedone
modeopen
monotonically_increasing_iddone
monthdone
monthsdone
months_betweenopen
nanvldone
next_dayopen
npopen
nth_valueopen
ntiledone
octet_lengthdone
overlayopen
overloadopen
pandas_udfopen
percent_rankdone
percentile_approxopen
pmodopen
posexplodedone
posexplode_outerdone
powdone
productdone
quarterdone
radiansdone
raise_erroropen
randdone
randndone
rankdone
regexp_extractopen
regexp_replaceopen
repeatopen
reversedone
rintdone
rounddone
row_numberdone
rpadopen
rtrimdone
schema_of_csvopen
schema_of_jsonopen
secdone
seconddone
sentencesopen
sequenceopen
session_windowopen
sha1done
sha2open
shiftLeftopen
shiftRightopen
shiftRightUnsignedopen
shiftleftopen
shiftrightopen
shiftrightunsignedopen
shuffledone
signumdone
sindone
sinhdone
sizedone
skewnessdone
sliceopen
sort_arrayopen
soundexdone
spark_partition_iddone
splitopen
sqrtdone
stddevdone
stddev_popdone
stddev_sampdone
structopen
substringopen
substring_indexopen
sumdone
sumDistinctopen
sum_distinctopen
sysopen
tandone
tanhdone
timestamp_secondsdone
toDegreesopen
toRadiansopen
to_csvopen
to_dateopen
to_jsonopen
to_stropen
to_timestampopen
to_utc_timestampopen
transformopen
transform_keysopen
transform_valuesopen
translateopen
trimdone
truncopen
try_remote_functionsopen
udfopen
unbase64done
unhexdone
unix_timestampopen
unwrap_udtopen
upperdone
var_popdone
var_sampdone
variancedone
warningsopen
weekofyeardone
whenopen
windowopen
window_timeopen
xxhash64done
yeardone
yearsdone
zip_withopen

Data Types

Data types are used for creating schemas and for casting columns to specific types

ColumnAPIComment
ArrayTypedone
BinaryTypedone
BooleanTypedone
ByteTypedone
DateTypedone
DecimalTypedone
DoubleTypedone
FloatTypedone
IntegerTypedone
LongTypedone
MapTypedone
NullTypedone
ShortTypedone
StringTypedone
CharTypedone
VarcharTypedone
StructFielddone
StructTypedone
TimestampTypedone
TimestampNTZTypedone
DayTimeIntervalTypedone
YearMonthIntervalTypedone

Literal Types

Create Spark literal types from these rust types. E.g. lit(1_i64) would be a LongType() in the schema.

An array can be made like lit([1_i16,2_i16,3_i16]) would result in an ArrayType(Short) since all the values of the slice can be translated into literal type.

Spark Literal TypeRust TypeStatus
Nullopen
Binary&[u8]done
Booleanbooldone
Byteopen
Shorti16done
Integeri32done
Longi64done
Floatf32done
Doublef64done
Decimalopen
String&str / Stringdone
Datechrono::NaiveDatedone
Timestampchrono::DateTime<Tz>done
TimestampNtzchrono::NaiveDateTimedone
CalendarIntervalopen
YearMonthIntervalopen
DayTimeIntervalopen
Arrayslice / Vecdone
MapCreate with the function create_mapdone
StructCreate with the function struct_col or named_structdone

Window & WindowSpec

For ease of use it's recommended to use Window to create the WindowSpec.

WindowAPIComment
currentRowdone
orderBydone
partitionBydone
rangeBetweendone
rowsBetweendone
unboundedFollowingdone
unboundedPrecedingdone
WindowSpec.orderBydone
WindowSpec.partitionBydone
WindowSpec.rangeBetweendone
WindowSpec.rowsBetweendone