replace_flow

Importante

Esta característica se encuentra en su versión beta.

El @dp.replace_flow decorador crea un flujo de REEMPLAZAR USANDO para una mesa de streaming en tu pipeline. En cada actualización, el flujo reemplaza todas las filas de la tabla de destino que coinciden con las replace_using columnas clave y deja todas las demás filas intactas. La función debe devolver un dataframe de streaming de Apache Spark. Véase Sustitución parcial de instantáneas con REEMPLAZAR flujos UTILIZADOS.

Úsalo @dp.replace_flow cuando tu fuente es una serie de instantáneas parciales codificadas por columna. Para definir la tabla objetivo y el flujo en una sola sentencia, en su lugar, pasa replace_using y sequence_by pasa a @dp.tabla.

Syntax

from pyspark import pipelines as dp

dp.create_streaming_table("<target-table-name>") # Required only if the target table doesn't exist.

@dp.replace_flow(
  target = "<target-table-name>",
  replace_using = ["<key-column>", "<key-column>"],
  sequence_by = "<sequence-column>",
  name = "<flow-name>", # optional, defaults to function name
  comment = "<comment>", # optional
  spark_conf = {"<key>" : "<value>", "<key>" : "<value>"}) # optional
def <function-name>():
  return (<streaming-query>)

Parámetros

Parámetro Tipo Descripción
function function Required. Función que devuelve un DataFrame de streaming de Apache Spark desde una consulta definida por el usuario.
target str Required. El nombre de la tabla de streaming que es el objetivo del flujo.
replace_using list Required. Las columnas clave que identifican qué filas objetivo reemplazar. Especifica al menos una columna. Las columnas clave no pueden repetirse, y el tipo de cada columna clave debe ser ordenable.
sequence_by str o Column Required. La columna que ordena las actualizaciones. Para cada clave, gana la secuencia más alta, y una fila de secuencia inferior nunca sobrescribe una que ya está en el objetivo.
name str Nombre del flujo. Si no se proporciona, el valor predeterminado es el nombre de la función.
comment str Descripción del flujo.
spark_conf dict Lista de configuraciones de Spark para la ejecución de esta consulta.

Examples

from pyspark import pipelines as dp

# Keep the latest row for each order from a stream of partial snapshots
dp.create_streaming_table("orders_current")

@dp.replace_flow(
  target = "orders_current",
  replace_using = ["order_id"],
  sequence_by = "updated_at"
)
def orders_flow():
  return spark.readStream.table("order_updates")

Utiliza más de una columna clave cuando un registro se identifica por una combinación de columnas:

from pyspark import pipelines as dp

dp.create_streaming_table("accounts_current")

@dp.replace_flow(
  target = "accounts_current",
  replace_using = ["region", "account_id"],
  sequence_by = "updated_at"
)
def accounts_flow():
  return spark.readStream.table("account_updates")

Limitaciones

  • Una tabla de flujo permite un solo REPLACE USING flujo y no puede combinarse REPLACE USING con otro tipo de flujo como un flujo de adición, un flujo automático CDC o un REPLACE WHERE flujo.
  • La consulta debe ser una consulta de streaming. @dp.replace_flow rechaza una fuente que no sea en streaming.
  • REPLACE USING los flujos requieren Databricks en tiempo de ejecución 18.2 y superiores.