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.
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 USINGflujo y no puede combinarseREPLACE USINGcon otro tipo de flujo como un flujo de adición, un flujo automático CDC o unREPLACE WHEREflujo. - La consulta debe ser una consulta de streaming.
@dp.replace_flowrechaza una fuente que no sea en streaming. -
REPLACE USINGlos flujos requieren Databricks en tiempo de ejecución 18.2 y superiores.