Portar una Aplicación a Apache Spark y Apache Flink

Las funciones con nombre hacen una consulta independiente de los operadores. Spark y Flink piden tres cambios más, porque difieren de PostgreSQL de tres maneras. Ya poseen algunos de los nombres que MobilityDB usa: ambos definen round, lower y length, y el analizador sintáctico de Flink rechaza set y unnest sin comillas; tal nombre lleva un prefijo que designa la clase o el tipo base de su argumento (“Los nombres que poseen los motores”). No llaman implícitamente ninguna función de entrada o salida de un tipo, de modo que un literal, una conversión y el texto de un valor se vuelven llamadas que la consulta nombra (“Lectura y escritura de valores”). Y cada motor llama a una función que devuelve filas a su manera (“Funciones que devuelven filas”). La tabla siguiente resume los cambios; una función que no cubre conserva su nombre.

Lo que contiene la consultaMobilityDBSpark y Flink
Un operador: su función con nombre, como la dan las secciones anteriorestrip && boxstboxOverlaps(trip, box)
Un nombre que el motor ya posee, sobre un valor cuyo tipo base tiene la operación: el tipo base delanteround(speed(trip), 2)floatRound(speed(trip), 2)
Un nombre que el motor ya posee, sobre otro valor conjunto, rango, conjunto de rangos o temporal: la clase delantelower(timeSpan(trip))spanLower(timeSpan(trip))
Un nombre que el motor ya posee, sobre un cuadro delimitador o un tipo base: el tipo delanteround(box, 2)stboxRound(box, 2)
Un literal con tipo: el constructor de texto de su tipotint '1@2001-01-01'tintFromText('1@2001-01-01')
El texto de un valor, que PostgreSQL escribe mediante la función de salida de su tipo: asTextSELECT speed(trip)SELECT asText(speed(trip))
Una conversión: la función que su entrada da a su lado, como tbool::tint es tint(tbool)flag::tinttint(flag)

Los nombres que poseen los motores

Un nombre que el motor posee lleva el prefijo de la clase o del tipo base de su argumento, y las funciones que van juntas llevan el mismo prefijo: spanLowerInc va con spanLower, geoCumulativeLength con geoLength. Spark y Flink llevan un único conjunto de nombres: cada uno registra tal función solo bajo su nombre con prefijo, de modo que una consulta escrita para uno se ejecuta en el otro, y round(1.23456, 2) sigue siendo el redondeo propio del motor de un número. Los nombres que cambian son los siguientes.

Nombre en MobilityDBClase del argumentoNombre en Spark y Flink
lower, upper, initcaptextset, ttexttextLower, textUpper, textInitcap
lower, upper, lowerInc, upperIncrango, conjunto de rangosspanLower, spanUpper, spanLowerInc, spanUpperInc, spansetLower, spansetUpper, spansetLowerInc, spansetUpperInc
lowerInc, upperInctemporaltemporalLowerInc, temporalUpperInc
roundfloat y sus tipos conjunto, rango, conjunto de rangos y temporalfloatRound
 geometría y geografía, sus conjuntos y tipos temporales, 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 y sus tipos conjunto, rango, conjunto de rangos y temporalfloatCeil, floatFloor, floatDegrees, floatRadians, floatCos, floatSin, floatTan, floatExp, floatLn, floatLog10
transformlos conjuntos y tipos temporales de geometría y geografía, trgeometrygeoTransform
transformPipelinegeometría y geografía, sus conjuntos y tipos temporales, trgeometrygeoTransformPipeline
 cbuffer, cbufferset, tcbuffer; pose, poseset, tpose; posechain, posechainset, tposechain; stbox; rastercbufferTransform, poseTransform, posechainTransform, stboxTransform, rasterTransform, y el …TransformPipeline de cada uno
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; lo mismo para posechain (posechainTranslate, …)
 trgeometrygeoTranslate, geoRotate, geoRotateX, geoRotateY, geoRotateZ
length, cumulativeLengthgeometría, geografía, tgeompoint, tgeogpoint, trgeometrygeoLength, geoCumulativeLength
 tnpointnpointLength, npointCumulativeLength
hash, hashExtendedconjunto, rango, conjunto de rangos, cuadro delimitador, temporal, tipo basesetHash, spanHash, spansetHash, tboxHash, stboxHash, temporalHash, cbufferHash, …, y sus …HashExtended
insert, update, mergetemporaltemporalInsert, temporalUpdate, temporalMerge
merge (el agregado)temporalmergeAgg
unnestconjunto, temporalsetUnnest, temporalUnnest
set (el constructor)cualquierasetMake

Lectura y escritura de valores

PostgreSQL lee y escribe un valor mediante las cuatro funciones que declara su tipo, y las llama implícitamente: un literal con tipo, una conversión desde texto y un COPY de texto pasan por la función de entrada, el texto que imprime un cliente por la función de salida, y un COPY binario o un controlador en modo binario por las funciones de recepción y envío. Spark y Flink no llaman ninguna de ellas: un valor viaja como sus bytes WKB, y cada dirección es una función que la consulta nombra, la inversa de otra, como muestra la tabla siguiente.

Lo que usa PostgreSQLDirecciónSpark y Flink
la función de entrada: un literal con tipo, una conversión desde textotexto a valorttypeFromText(text), settypeFromText(text), boxFromText(text), …, una por tipo, con el nombre del tipo como prefijo; con un SRID también tspatialFromEWKT(text)
la función de salida: el texto que imprime un clientevalor a textoasText(value); con un SRID también asEWKT(value)
la función de recepción: COPY binario, un controlador binariobinario a valorttypeFromBinary(bytea), ttypeFromHexWKB(text), y lo mismo para todo otro tipo; con un SRID también tspatialFromEWKB(bytea), tspatialFromHexEWKB(text)
la función de envíovalor a binarioasBinary(value), asHexWKB(value); con un SRID también asEWKB(value), asHexEWKB(value)
ningunaMF-JSONttypeFromMFJSON(text), asMFJSON(ttype)

Todo tipo tiene cada par que puede expresar, las formas E cuando sus formas de texto y binaria omiten el SRID que lleva y MF-JSON cuando es temporal, incluidos los tipos que PostgreSQL lee y escribe solo mediante sus cuatro funciones: los tipos de conjunto, de rango y de conjunto de rangos de enteros, enteros grandes, flotantes, fechas y marcas de tiempo, textset y jsonbset, tbox y stbox, tbool, tint, tbigint, tfloat y ttext, y quadbin, s2cell, nsegment y tpcbox. Un tipo cuyo texto es su WKB hexadecimal, raquet, lee y escribe su texto mediante raquetFromHexWKB y asHexWKB, como un ráster de PostGIS mediante ST_RastFromHexWKB y ST_AsHexWKB. Una caja, stbox o tpcbox, escribe su SRID en su texto y en su WKB, por lo que no necesita formas E. Toda salida WKB acepta el orden de bytes endian, y toda salida de texto de un valor con números flotantes o coordenadas la precisión maxdecimaldigits. Una columna de bytes WKB, por ejemplo una columna TemporalParquet, entra así en todos los motores por el constructor binario de su tipo, como en tgeompointFromBinary(col).

Funciones que devuelven filas

Una función que devuelve un conjunto de filas, como timeSplit o unnest, se llama como función de tabla en la cláusula FROM en PostgreSQL. En Flink y en Spark la misma función, unnest bajo los nombres setUnnest y temporalUnnest, devuelve un arreglo, un elemento por fila, un arreglo de registros cuando devuelve varias columnas; Flink lo despliega en filas con CROSS JOIN UNNEST, y Spark con LATERAL VIEW explode, o inline para un arreglo de registros.

Ejemplo

La consulta

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

se escribe, en Spark y Flink,

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

donde round del número resultante es el del propio motor.