Nota:
El acceso a esta página requiere autorización. Puede intentar iniciar sesión o cambiar directorios.
El acceso a esta página requiere autorización. Puede intentar cambiar los directorios.
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:
|
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_checkpointsolo se soporta en append_flow.