Nota
O acesso a esta página requer autorização. Pode tentar iniciar sessão ou alterar os diretórios.
O acesso a esta página requer autorização. Pode tentar alterar os diretórios.
Use o @dp.update_flow decorador para criar um fluxo de atualizações. Os fluxos de atualização escrevem para os sinks usando o modo de saída de atualização, emitindo apenas as linhas que mudaram em cada lote. Ao contrário dos fluxos anexos, suportam agregações com estado sem necessidade de marca de água.
Os fluxos de atualização só podem atingir os sumidouros. As tabelas Delta não são suportadas.
Syntax
from pyspark import pipelines as dp
dp.create_sink("<sink-name>", "<format>", {"<key>": "<value>"})
@dp.update_flow(
target = "<sink-name>",
name = "<flow-name>", # optional, defaults to function name
depends_on = "<flow-name>", # optional, Public Preview
spark_conf = {"<key>" : "<value>", "<key>" : "<value>"}, # optional
comment = "<comment>") # optional
def <function-name>():
return (<streaming-query>)
Parameters
| Parâmetro | Tipo | Description |
|---|---|---|
| função | function |
Required. Uma função que retorna um DataFrame de streaming Apache Spark de uma consulta definida pelo usuário. |
target |
str |
Required. O nome do sumidouro onde este fluxo escreve. |
name |
str |
O nome do fluxo. Se não for fornecido, o padrão será o nome da função. |
depends_on |
str ou list |
Pré-visualização pública. Um ou mais nomes de fluxo que devem ser concluídos com sucesso antes de este fluxo começar. Aceita um único nome de fluxo ou uma lista de nomes. Isto ordena apenas a execução do fluxo; não altera a forma como o fluxo corre. Veja execução de fluxo de pipeline de ordens com depends_on. |
comment |
str |
Uma descrição para o fluxo. |
spark_conf |
dict |
Um ditado das configurações do Spark para a execução desta consulta. Estas configurações sobrepõem as conferências definidas para o destino, pipeline ou cluster. |
Exemplos
Agregação a um sumidouro de Kafka
Escreva resultados de agregação com estado num sumidouro de Kafka:
from pyspark import pipelines as dp
from pyspark.sql.functions import col
dp.create_sink("event_counts_sink", "kafka", {
"kafka.bootstrap.servers": broker_address,
"topic": output_topic,
})
@dp.update_flow(
name="event_counts_flow",
target="event_counts_sink",
)
def event_counts():
return (
spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", broker_address)
.option("subscribe", input_topic)
.load()
.selectExpr("CAST(key AS STRING) AS event_type")
.groupBy(col("event_type"))
.count()
)
Modo em tempo real
Importante
O modo em tempo real está em Pré-visualização Pública.
Use spark_conf para configurar um fluxo de atualização para o modo em tempo real:
from pyspark import pipelines as dp
dp.create_sink("my_kafka_sink", "kafka", {
"kafka.bootstrap.servers": broker_address,
"topic": output_topic,
})
@dp.update_flow(
name="my_rtm_flow",
target="my_kafka_sink",
spark_conf={
"pipelines.trigger": "RealTime",
"pipelines.trigger.interval": "5 minutes",
}
)
def my_real_time_flow():
return (
spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", broker_address)
.option("subscribe", input_topic)
.load()
)
Limitações
- Os sumidouros de tabela Delta não são suportados como alvos para fluxos de atualização.