Ir al contenido
SparquetSparquet

Conectores

Un conector se elige mediante el campo format de un input, un output o un origen de join/union. El mismo registry sirve para los tres, y es extensible — consulta Agregar el tuyo.

Formato Lectura Escritura Notas
parquet Nativo de Spark
delta Upserts con MERGE, time travel
iceberg MERGE INTO
csv header e inferSchema toman true por defecto; comillas RFC 4180
txt texto plano, una sola columna value
view vistas temporales de Spark, auto-cacheadas
kafka lectura y publicación batch; Amazon MSK
postgresql JDBC
mysql JDBC
mariadb JDBC
sqlserver JDBC
oracle JDBC (service name)
bigquery Google BigQuery
snowflake Snowflake
redshift Amazon Redshift, staging en S3
mongodb MongoDB
documentdb Amazon DocumentDB (protocolo Mongo)
dynamodb Amazon DynamoDB, la escritura es upsert
cassandra Cassandra / ScyllaDB, la escritura es upsert
elasticsearch Elasticsearch / OpenSearch
{ "format": "parquet", "path": "/data/curated/orders", "mode": "overwrite", "partition_by": ["country"] }

Cualquier opción Parquet de Spark se pasa a través de options. path es un directorio de archivos part, no un archivo único.

{ "format": "csv", "path": "/data/landing/orders", "options": { "sep": ";", "encoding": "UTF-8" } }

Valores por defecto que aplica el framework: header: true e inferSchema: true en lectura, header: true en escritura, encoding: UTF-8 y escape: " en ambas. inferSchema cuesta una pasada extra sobre los datos — define un cast explícito y desactívalo para feeds grandes y estables.

Las comillas siguen el RFC 4180: una comilla dentro de un campo se escribe duplicada (""), no escapada con barra invertida (\") como hace Spark por defecto. Spark relee su propio dialecto, pero el csv de Python, pandas y Excel no — eso rompía el validations.report, cuya columna rule_params es JSON lleno de comillas, precisamente en las herramientas donde se analiza. Para leer un archivo escrito en el dialecto antiguo, decláralo: "options": { "escape": "\\" }.

// read a table, or a path, or a past version
{ "format": "delta", "path": "catalog.schema.orders" }
{ "format": "delta", "path": "/mnt/raw/orders", "options": { "versionAsOf": "5" } }
{ "format": "delta", "path": "catalog.schema.orders", "options": { "timestampAsOf": "2026-01-01T00:00:00Z" } }
// upsert
{ "format": "delta", "path": "catalog.schema.orders", "mode": "merge",
"options": { "merge_keys": ["order_id"], "merge_condition": "T.deleted = false" } }

path se trata como nombre de tabla cuando contiene un punto y no empieza con / ni con un esquema de almacenamiento; de lo contrario, como una ruta. En merge_condition, T es la tabla destino y S el DataFrame entrante.

Fuera de Databricks, instala el paquete OSS: pip install "sparquet[delta]".

{ "format": "iceberg", "path": "catalog.db.orders", "mode": "merge", "options": { "merge_keys": ["id"] } }
{ "format": "txt", "path": "/data/exports/lines", "mode": "overwrite" }

Lee hacia una sola columna value, y escribe un DataFrame que debe tener exactamente una columna de tipo string — proyéctala antes con select.

{ "format": "view", "path": "orders_staging", "mode": "overwrite" }

Registra una vista temporal de Spark, cacheada por defecto, de modo que un pipeline posterior en la misma sesión pueda leerla como {"format": "view", "path": "orders_staging"}. Este es el mecanismo de staging detrás de los jobs multi-pipeline — uno escribe la vista, el siguiente la consume.

options.scope elige el tiempo de vida: session (por defecto) es una vista temporal normal, visible solo para la sesión actual; global usa createOrReplaceGlobalTempView, visible para cada sesión de la misma aplicación Spark (léela como global_temp.<name>, o con scope: "global" en el reader).

Lectura batch y publicación batch. path es el topic. Requiere spark-sql-kafka-0-10 en el classpath.

// read the whole topic
{
"format": "kafka",
"path": "orders-events",
"options": { "bootstrap_servers": "broker1:9092,broker2:9092", "startingOffsets": "earliest" }
}
// publish
{
"format": "kafka",
"path": "orders-events",
"mode": "append",
"transformations": [
{ "type": "with_column", "column": "value", "expression": "to_json(payload)" }
],
"options": {
"bootstrap_servers": "broker1:9092,broker2:9092",
"value_column": "value",
"key_column": "order_id"
}
}

En lectura, el DataFrame regresa con el esquema nativo de Kafka (key/value binarios, topic, partition, offset, timestamp, timestampType) — normalmente seguido de un CAST(value AS STRING). Los valores por defecto del batch son startingOffsets: earliest y endingOffsets: latest, así que una ejecución lee el topic completo; sobrescríbelos, o pasa assign / subscribePattern en lugar del subscribe por defecto.

En escritura, la value_column configurada (por defecto payload) y la key_column (por defecto header, null para omitirla) se renombran a value y key, y toda columna fuera de {key, value, topic, partition, timestamp, headers} se descarta antes de la escritura.

bootstrap_servers es un alias amigable de kafka.bootstrap.servers. Para Amazon MSK con autenticación IAM, pasa kafka.security.protocol=SASL_SSL y kafka.sasl.mechanism=AWS_MSK_IAM en options (además del JAR aws-msk-iam-auth en el classpath).

Nativo, sin JAR extra. path es un directorio; el valor por defecto es un objeto por línea (JSON Lines). Define multiLine: "true" para leer un documento por archivo.

{ "format": "json", "path": "/data/landing/events", "options": { "multiLine": "true" } }

Formato columnar nativo (como Parquet), sin JAR extra. Las opciones incluyen compression (por defecto zlib) y mergeSchema.

{ "format": "orc", "path": "/data/curated/orders", "mode": "overwrite", "partition_by": ["dt_ref"] }

Orientado a filas; requiere el paquete org.apache.spark:spark-avro en el classpath. Opciones: avroSchema, compression (snappy/deflate/bzip2/xz), recordName.

{ "format": "avro", "path": "/data/raw/events" }

Requiere el paquete com.databricks:spark-xml (registra el formato xml). rowTag es obligatorio en lectura y escritura; rootTag (por defecto rows) nombra la raíz de escritura.

{ "format": "xml", "path": "/data/raw/catalog", "options": { "rowTag": "book" } }

Solo lectura (binaryFile). Carga archivos completos en las columnas path, modificationTime, length, content (binario) — imágenes, PDFs, blobs. No hay writer binario; persiste content vía parquet/delta.

{ "format": "binary", "path": "/data/raw/docs", "options": { "pathGlobFilter": "*.pdf" } }

Formato de tabla lakehouse con upserts. Requiere el JAR hudi-spark-bundle y las extensiones de sesión de Hudi. path es la ruta base de la tabla; el particionado y el upsert se gobiernan mediante opciones hoodie.* (el partition_by del framework no se usa).

{
"format": "hudi",
"path": "s3a://lake/hudi/orders",
"mode": "append",
"options": {
"hoodie.table.name": "orders",
"hoodie.datasource.write.recordkey.field": "order_id",
"hoodie.datasource.write.precombine.field": "updated_at",
"hoodie.datasource.write.operation": "upsert"
}
}

postgresql, mysql, mariadb, sqlserver y oracle comparten una base JDBC común con dialectos por base de datos que completan la clase del driver, el puerto por defecto y la forma de la URL. path es el nombre de la tabla (dbtable).

// read a table
{
"format": "postgresql",
"path": "public.orders",
"options": {
"host": "db.internal", "port": 5432, "database": "sales",
"user": "reader", "password": ""
}
}
// or give the full URL and a pushdown query
{
"format": "sqlserver",
"path": "dbo.orders",
"options": {
"url": "jdbc:sqlserver://db:1433;databaseName=sales",
"query": "SELECT id, total FROM dbo.orders WHERE total > 0",
"user": "reader", "password": ""
}
}

Opciones de conexión: url (URL JDBC completa, tiene prioridad) o host + database (+ opcionales port, con valor por defecto por base de datos) para construirla; driver (con valor por defecto por base de datos); user / password. Las lecturas también aceptan query (un SELECT usado en lugar de dbtable), dbtable (sobrescribe path), partitionColumn / lowerBound / upperBound / numPartitions para lecturas paralelas, y fetchsize. Las escrituras soportan append y overwrite (truncate: "true" hace un TRUNCATE en lugar de DROP/CREATE), además de batchsize, isolationLevel, createTableColumnTypes / createTableOptions. JDBC no tiene merge — usa append/overwrite.

Spark puede entregar partes de la consulta a la base de datos en vez de arrastrar filas por la red para filtrarlas en la JVM. Las cinco palancas tienen valor por defecto true en el camino de lectura de Spark 4 y se pasan directamente:

Opción Qué va a la base de datos
pushDownPredicate la cláusula WHERE
pushDownAggregate sum/count/min/max/avg y el GROUP BY
pushDownLimit el LIMIT
pushDownOffset el OFFSET
pushDownTableSample TABLESAMPLE
{ "format": "postgresql", "path": "public.orders",
"options": { "host": "db.internal", "database": "shop",
"pushDownAggregate": "true", "pushDownLimit": "true" } }

La restricción que importa: solo se aplican cuando la lectura es de una tabla. Con query el SELECT que escribiste ya es el corte, así que Spark lo envuelve como subconsulta y no empuja nada más dentro de él. El pushdown de agregación también exige que la agregación sea lo único encima del scan — un filter sobre una columna calculada, un join o una UDF en medio y la base de datos recibe un SELECT simple. Activar una palanca es una petición, no un hecho: { "type": "debug", "actions": ["pushdown"] } dice qué empujó realmente el plan. Consulta debug.

mariadb construye una URL jdbc:mysql:// con el driver de MySQL, a propósito. Spark 4.1 no tiene dialecto MariaDB y su MySQLDialect solo coincide con URLs que empiezan por jdbc:mysql; con jdbc:mariadb:// el conector cae al dialecto por defecto, que cita identificadores con " — y MariaDB lo rechaza tanto en lecturas como en escrituras, ya que el SELECT que Spark construye cita las columnas del mismo modo:

You have an error in your SQL syntax ... near '"id" INTEGER'

MariaDB habla el protocolo de MySQL, así que Connector/J alcanza un servidor MariaDB y trae el dialecto consigo: citado con backtick, mapeo de tipos de MySQL y SQL correcto en el pushdown de agregación y de LIMIT. Para usar el driver de MariaDB de todos modos, pasa url y driver explícitamente junto con sessionVariables:

{ "format": "mariadb", "path": "orders",
"options": { "url": "jdbc:mariadb://db.internal:3306/shop?sessionVariables=sql_mode='ANSI_QUOTES'",
"driver": "org.mariadb.jdbc.Driver", "user": "app", "password": "..." } }

El precio es que "..." deja de ser un literal de cadena en esa sesión, lo que importa a quien use query. Una URL jdbc:mariadb:// explícita sin ANSI_QUOTES emite una advertencia.

Para Oracle, database es el service name (jdbc:oracle:thin:@//host:port/service).

path es project.dataset.table (o dataset.table con un proyecto por defecto). Usa el spark-bigquery-connector.

{ "format": "bigquery", "path": "my-proj.sales.orders",
"options": { "parentProject": "billing-proj", "credentialsFile": "/secrets/sa.json" } }

Opciones de lectura: query (requiere viewsEnabled: "true"), parentProject (facturación), credentialsFile / credentials, filter, maxParallelism. Escritura (overwrite / append): temporaryGcsBucket (staging indirecto, por defecto), writeMethod (indirect o direct), partitionField / partitionType / clusteredFields.

path es la tabla (dbtable). Usa spark-snowflake + snowflake-jdbc.

{ "format": "snowflake", "path": "ANALYTICS.PUBLIC.ORDERS",
"options": {
"sfUrl": "myorg-acct.snowflakecomputing.com",
"sfUser": "loader", "sfPassword": "",
"sfDatabase": "ANALYTICS", "sfSchema": "PUBLIC", "sfWarehouse": "LOAD_WH"
} }

La conexión es la familia sfXxx (sfUrl, sfUser / sfPassword o pem_private_key, sfDatabase, sfSchema, sfWarehouse, sfRole). Las lecturas aceptan query; las escrituras soportan overwrite / append.

path es la tabla (dbtable). Usa el conector comunitario spark-redshift, que hace staging a través de S3.

{ "format": "redshift", "path": "public.orders",
"options": {
"url": "jdbc:redshift://cluster:5439/sales",
"tempdir": "s3://my-bucket/redshift-staging/",
"aws_iam_role": "arn:aws:iam::…:role/redshift-copy"
} }

url y tempdir (un prefijo S3) son obligatorios. La autenticación es user / password o aws_iam_role; forward_spark_s3_credentials: "true" reutiliza las credenciales S3 de la sesión. Las escrituras soportan overwrite / append y diststyle / distkey / sortkeyspec / tempformat.

path es la colección. El mismo Mongo Spark Connector sirve para Amazon DocumentDB — apunta connection.uri al clúster de DocumentDB con ?tls=true&retryWrites=false.

{ "format": "mongodb", "path": "orders",
"options": { "connection.uri": "mongodb://user:pass@host:27017/", "database": "sales" } }

Opciones: connection.uri (obligatoria), database, collection (sobrescribe path), aggregation.pipeline (pushdown de lectura). Las escrituras soportan overwrite / append, además de operationType (insert / replace / update), idFieldList, ordered.

path es la tabla (tableName). Usa spark-dynamodb.

{ "format": "dynamodb", "path": "orders", "options": { "region": "us-east-1" } }

Opciones: region, roleArn, endpoint (p. ej. DynamoDB local), throughput / targetCapacity, readPartitions / stronglyConsistentReads. Las escrituras son siempre un PutItem por fila — un upsert por clave primaria (append); DynamoDB no tiene overwrite de tabla.

path es keyspace.table (o solo la tabla, con keyspace en options). Usa el spark-cassandra-connector.

{ "format": "cassandra", "path": "sales.orders",
"options": { "spark.cassandra.connection.host": "node1,node2" } }

Opciones: keyspace / table (sobrescriben path), spark.cassandra.connection.host / .port, spark.cassandra.auth.username / .password. Las escrituras son append (INSERT/upsert por clave); la tabla debe existir de antemano.

path es el índice (resource). Usa el conector elasticsearch-hadoop — prefijo de opciones es.*.

{ "format": "elasticsearch", "path": "orders",
"options": { "es.nodes": "es.internal", "es.port": "9200" } }

Opciones: es.nodes / es.port, es.net.http.auth.user / .pass, es.nodes.wan.only (cloud/proxy), es.query (DSL de lectura), es.mapping.id (usar una columna como _id), es.write.operation (index / create / update / upsert). Las escrituras soportan append / overwrite.

Un formato separado con su propio conector (opensearch-hadoop) — prefijo de opciones opensearch.*, no es.*. Por lo demás, la misma forma que Elasticsearch.

{ "format": "opensearch", "path": "orders",
"options": { "opensearch.nodes": "os.internal", "opensearch.port": "9200" } }

Las opciones reflejan las de ES con el prefijo opensearch.: opensearch.nodes / opensearch.port, opensearch.net.http.auth.user / .pass, opensearch.nodes.wan.only, opensearch.query, opensearch.mapping.id, opensearch.write.operation.

En Spark 4 el conector es org.opensearch.client:opensearch-spark-40_2.13:2.0.0; en Spark 3.x es opensearch-spark-30_2.12. Esta es también la vía soportada hacia un servidor Elasticsearch — consulta la advertencia de arriba.

Cada conector de arriba se distribuye como un paquete Spark que debe estar en el classpath. Decláralo por pipeline para que el job lleve su propia dependencia:

{
"spark": {
"configs": {
"spark.jars.packages": "org.postgresql:postgresql:42.7.3,org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.1"
}
}
}

En una plataforma administrada también puedes agregar el paquete a las librerías del clúster. El framework nunca los empaqueta — un JAR faltante se manifiesta como un ClassNotFoundException / Failed to find data source en el momento de lectura o escritura.

Registra un reader o writer una vez y todo pipeline del proceso podrá usarlo:

fw.register_reader("my_format", MyReader) # class MyReader(BaseReader)
fw.register_writer("my_format", MyWriter) # class MyWriter(BaseWriter)

Consulta Extender. Studio preserva los formatos desconocidos a través de la importación y exportación, así que un conector personalizado no rompe el canvas.