append_flow

El @dp.append_flow decorador crea flujos de anexión o reposición para las tablas de canalización. La función debe devolver un dataframe de streaming de Apache Spark. Consulte Carga y procesamiento de datos de forma incremental con flujos de canalización de Lakeflow.

Los flujos de adición pueden dirigirse a tablas de streaming, tablas gestionadas o sumideros.

Syntax

from pyspark import pipelines as dp

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

@dp.append_flow(
  target = "<target-table-name>",
  name = "<flow-name>", # optional, defaults to function name
  once = False, # optional
  depends_on = "<flow-name>", # optional, Public Preview
  spark_conf = {"<key>" : "<value", "<key" : "<value>"}, # optional
  comment = "<comment>", # optional
  import_checkpoint = "<checkpoint-path>") # optional
def <function-name>():
  return (<streaming-query>) #

Parámetros

Parámetro Tipo Description
function function Obligatorio. Función que devuelve un DataFrame de streaming de Apache Spark desde una consulta definida por el usuario.
target str Obligatorio. Nombre de la tabla o receptor que es el destino del flujo de anexión.
name str Nombre del flujo. Si no se proporciona, el valor predeterminado es el nombre de la función.
once bool Opcionalmente, defina el flujo como un flujo de un solo uso, como un reposición. El uso de once=True cambia el flujo de dos maneras:
  • El valor de retorno. streaming-query. debe ser un dataframe por lotes en este caso, no un dataframe de streaming.
  • El flujo se ejecuta una vez de forma predeterminada. Si la canalización se actualiza por completo, el flujo ONCE se ejecuta nuevamente para recrear los datos.
depends_on str o list Versión preliminar pública. Uno o más nombres de flujo deben completarse con éxito antes de que este flujo comience. Acepta un solo nombre de flujo o una lista de nombres. Esto ordena solo la ejecución del flujo; no cambia cómo se ejecuta el flujo. Consulta la ejecución de flujo de la tubería de órdenes con depends_on.
comment str Descripción del flujo.
spark_conf dict Lista de configuraciones de Spark para la ejecución de esta consulta
import_checkpoint str La ruta hacia un punto de control de Streaming Estructurado existente para importarlo al flujo, de modo que un flujo migrado se reanude desde su último offset comprometido en lugar de reprocesar la fuente. Importar un punto de control está en Beta. Consulta Migrar un punto de control de Streaming Estructurado.

Examples

from pyspark import pipelines as dp

# Create a sink for an external Delta table
dp.create_sink("my_sink", "delta", {"path": "/tmp/delta_sink"})

# Add an append flow to an external Delta table
@dp.append_flow(name = "flow", target = "my_sink")
def flowFunc():
  return <streaming-query>

# Add a backfill
@dp.append_flow(name = "backfill", target = "my_sink", once = True)
def backfillFlowFunc():
    return (
      spark.read
      .format("json")
      .load("/path/to/backfill/")
    )

# Create a Kafka sink
dp.create_sink(
  "my_kafka_sink",
  "kafka",
  {
    "kafka.bootstrap.servers": "host:port",
    "topic": "my_topic"
  }
)

# Add an append flow to a Kafka sink
@dp.append_flow(name = "flow", target = "my_kafka_sink")
def myFlow():
  return read_stream("xxx").select(F.to_json(F.struct("*")).alias("value"))

Migrar un punto de control de Streaming Estructurado

Importante

Importar un punto de control está en Beta.

Se utiliza import_checkpoint para migrar una carga de trabajo de Streaming Estructurado existente a una tubería sin reprocesar el código fuente. Configúralo como la checkpointLocation consulta de Streaming Estructurado que uses, que puede ser un almacenamiento en la nube, un volumen de Unity Catalog o una ruta DBFS. En la primera actualización de la canalización, el flujo clona ese punto de control en el almacenamiento gestionado de la canalización. El flujo reanuda entonces desde el último desplazamiento comprometido con su estado (como agregaciones, claves de deduplicación y marcas de agua) intacto. Las actualizaciones posteriores de la tubería utilizan el punto de control clonado del flujo; el punto de control original no se modifica.

El flujo debe dirigirse a una tabla gestionada creada con create_table o un sumidero.

Detiene la consulta original de Structured Streaming antes de ejecutar la pipeline. La consulta original de Structured Streaming puede reutilizarse tras la importación, pero necesitas gestionar su estado de punto de control y asegurarte de que la pipeline y la consulta no escriban en la misma tabla al mismo tiempo, lo que puede producir datos duplicados.

Recrea la consulta de Structured Streaming como un flujo de pipeline que escribe en una nueva tabla e importa su punto de control:

from pyspark import pipelines as dp

# Create a new managed table for the pipeline
dp.create_table("target_table")

# Continue from the imported checkpoint instead of reprocessing the source.
@dp.append_flow(
  target = "target_table",
  import_checkpoint = "/Volumes/my_catalog/my_schema/checkpoints/my_stream",
)
def migrate():
  # The same source your original query read from.
  return spark.readStream.table("source_table")

El punto de control se importa solo una vez, en la primera actualización del pipeline; las actualizaciones posteriores ignoran import_checkpoint. Una actualización completa no reimporta el punto de control; comienza desde un nuevo punto vacío y reprocesa el origen. Para importar un punto de control diferente, usa un nombre de flujo que no se haya usado antes para la tabla objetivo; reutilizar un nombre de flujo existente salta la importación.

Limitaciones

  • No se soporta importar un punto de control en una tabla que ya existe (por ejemplo, el objetivo original de consulta de Structured Streaming). Apunta a una nueva tabla que cree la pipeline, o a un sumidero.
  • import_checkpoint solo se soporta en append_flow.