Porting an Application to Apache Spark and Apache Flink

The named functions make a query independent of operators. Spark and Flink ask for three more changes, because they differ from PostgreSQL in three ways. They already own some of the names MobilityDB uses: both define round, lower and length, and Flink's parser refuses set and unnest unquoted; such a name takes a prefix naming the class or base type of its argument (the section called “Names the Engines Own”). They call no input or output function of a type implicitly, so a literal, a cast and the text of a value become calls the query names (the section called “Reading and Writing Values”). And each engine calls a function returning rows its own way (the section called “Functions Returning Rows”). The following table summarizes the changes; a function it does not cover keeps its name.

What the query holdsMobilityDBSpark and Flink
An operator: its named function, as the sections above give ittrip && boxstboxOverlaps(trip, box)
A name the engine owns, over a value whose base type has the operation: the base type in frontround(speed(trip), 2)floatRound(speed(trip), 2)
A name the engine owns, over another set, span, span set or temporal value: the class in frontlower(timeSpan(trip))spanLower(timeSpan(trip))
A name the engine owns, over a box or a base type: the type in frontround(box, 2)stboxRound(box, 2)
A typed literal: the text constructor of its typetint '1@2001-01-01'tintFromText('1@2001-01-01')
The text of a value, which PostgreSQL writes through the output function of its type: asTextSELECT speed(trip)SELECT asText(speed(trip))
A cast: the function its entry gives beside it, as tbool::tint is tint(tbool)flag::tinttint(flag)

Names the Engines Own

A name the engine owns takes the prefix of the class or base type of its argument, and functions that belong together take the same prefix: spanLowerInc goes with spanLower, geoCumulativeLength with geoLength. Spark and Flink carry one set of names: each registers such a function under its prefixed name alone, so a query written for one runs on the other, and round(1.23456, 2) remains the engine's own rounding of a number. The names that change are the following.

MobilityDB nameClass of the argumentSpark and Flink name
lower, upper, initcaptextset, ttexttextLower, textUpper, textInitcap
lower, upper, lowerInc, upperIncspan, span setspanLower, spanUpper, spanLowerInc, spanUpperInc, spansetLower, spansetUpper, spansetLowerInc, spansetUpperInc
lowerInc, upperInctemporaltemporalLowerInc, temporalUpperInc
roundfloat and its set, span, span set and temporal typesfloatRound
 geometry and geography, their sets and temporal types, trgeometrygeoRound
 cbuffer, cbufferset, tcbuffer; npoint, npointset, tnpoint; nsegment; pose, poseset, tpose; posechain, posechainset, tposechain; tbox; stbox; tpcboxcbufferRound, npointRound, nsegmentRound, poseRound, posechainRound, tboxRound, stboxRound, tpcboxRound
abstint, tbigint, tfloatintAbs, bigintAbs, floatAbs
ceil, floor, degrees, radians, cos, sin, tan, exp, ln, log10float and its set, span, span set and temporal typesfloatCeil, floatFloor, floatDegrees, floatRadians, floatCos, floatSin, floatTan, floatExp, floatLn, floatLog10
transformthe sets and temporal types of geometry and geography, trgeometrygeoTransform
transformPipelinegeometry and geography, their sets and temporal types, trgeometrygeoTransformPipeline
 cbuffer, cbufferset, tcbuffer; pose, poseset, tpose; posechain, posechainset, tposechain; stbox; rastercbufferTransform, poseTransform, posechainTransform, stboxTransform, rasterTransform, and the …TransformPipeline of each
translate, affine, rotate, rotateX, rotateY, rotateZ, scale, transscalegeometry, tgeompoint, tgeometrygeoTranslate, geoAffine, geoRotate, geoRotateX, geoRotateY, geoRotateZ, geoScale, geoTransscale
translate, rotate, rotateX, rotateY, rotateZtcbuffer; tpose; tposechaincbufferTranslate, cbufferRotate, cbufferRotateZ; poseTranslate, poseRotate, poseRotateX, poseRotateY, poseRotateZ; the same for posechain (posechainTranslate, …)
 trgeometrygeoTranslate, geoRotate, geoRotateX, geoRotateY, geoRotateZ
length, cumulativeLengthgeometry, geography, tgeompoint, tgeogpoint, trgeometrygeoLength, geoCumulativeLength
 tnpointnpointLength, npointCumulativeLength
hash, hashExtendedset, span, span set, box, temporal, base typesetHash, spanHash, spansetHash, tboxHash, stboxHash, temporalHash, cbufferHash, …, and their …HashExtended
insert, update, mergetemporaltemporalInsert, temporalUpdate, temporalMerge
merge (the aggregate)temporalmergeAgg
unnestset, temporalsetUnnest, temporalUnnest
set (the constructor)anysetMake

Reading and Writing Values

PostgreSQL reads and writes a value through the four functions its type declares, and calls them implicitly: a typed literal, a cast from text and a text COPY go through the input function, the text a client prints through the output function, and a binary COPY or a driver in binary mode through the receive and send functions. Spark and Flink call none of them: a value travels as its WKB bytes, and each direction is a function the query names, the inverse of another, as the following table shows.

What PostgreSQL usesDirectionSpark and Flink
the input function: a typed literal, a cast from texttext to valuettypeFromText(text), settypeFromText(text), boxFromText(text), …, one per type, prefixed by its name; with an SRID also tspatialFromEWKT(text)
the output function: the text a client printsvalue to textasText(value); with an SRID also asEWKT(value)
the receive function: binary COPY, a binary driverbinary to valuettypeFromBinary(bytea), ttypeFromHexWKB(text), and likewise for every other type; with an SRID also tspatialFromEWKB(bytea), tspatialFromHexEWKB(text)
the send functionvalue to binaryasBinary(value), asHexWKB(value); with an SRID also asEWKB(value), asHexEWKB(value)
noneMF-JSONttypeFromMFJSON(text), asMFJSON(ttype)

Every type has each pair it can state, the E forms when its text and binary forms omit the SRID it carries and MF-JSON when it is temporal, the types PostgreSQL reads and writes only through its four functions included: the set, span and span set types of integers, big integers, floats, dates and timestamps, textset and jsonbset, tbox and stbox, tbool, tint, tbigint, tfloat and ttext, and quadbin, s2cell, nsegment and tpcbox. A type whose text is its hex WKB, raquet, reads and writes its text through raquetFromHexWKB and asHexWKB, as a PostGIS raster through ST_RastFromHexWKB and ST_AsHexWKB. A box, stbox or tpcbox, writes its SRID in its text and its WKB, so it needs no E form. Every WKB output takes the byte order endian, and every text output of a value carrying floats or coordinates the precision maxdecimaldigits. A column of WKB bytes, a TemporalParquet column for instance, thus enters every engine through the binary constructor of its type, as in tgeompointFromBinary(col).

Functions Returning Rows

A function returning a set of rows, such as timeSplit or unnest, is called as a table function in the FROM clause in PostgreSQL. In Flink and in Spark the same function, unnest under the names setUnnest and temporalUnnest, returns an array, one element per row, an array of records when it returns several columns; Flink unfolds it into rows with CROSS JOIN UNNEST, and Spark with LATERAL VIEW explode, or inline for an array of records.

Example

The query

SELECT vehicle, round(length(trip)::numeric, 1), lower(timeSpan(trip))
FROM Trips
WHERE trip && stbox 'SRID=3857;STBOX X((0,0),(1000,1000))';

is written, in Spark and Flink,

SELECT vehicle, round(geoLength(trip), 1), spanLower(timeSpan(trip))
FROM Trips
WHERE stboxOverlaps(trip, stboxFromText('SRID=3857;STBOX X((0,0),(1000,1000))'));

where round of the resulting number is the engine's own.