@dp.append_flow 데코레이터는 파이프라인 테이블에 대한 추가 흐름 또는 백필을 생성합니다. 함수는 Apache Spark 스트리밍 데이터 프레임을 반환해야 합니다.
Lakeflow 파이프라인 흐름을 사용하여 데이터 증분 로드 및 처리를 참조하세요.
Append 플로우는 스트리밍 테이블, 관리 테이블, 또는 싱크를 대상으로 할 수 있습니다.
문법
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>) #
매개 변수
| 매개 변수 | 유형 | Description |
|---|---|---|
| 기능 | function |
필수 사항입니다. 사용자 정의 쿼리에서 Apache Spark 스트리밍 DataFrame을 반환하는 함수입니다. |
target |
str |
필수 사항입니다. 추가 흐름의 대상인 테이블 또는 싱크의 이름입니다. |
name |
str |
흐름 이름입니다. 제공되지 않으면 기본적으로 함수 이름이 지정됩니다. |
once |
bool |
필요에 따라 흐름을 백필과 같은 일회성 흐름으로 정의합니다.
once=True 다음 두 가지 방법으로 흐름을 변경합니다.
|
depends_on |
str 또는 list |
공개 미리 보기. 이 흐름이 시작되기 전에 반드시 성공적으로 완료되어야 하는 하나 이상의 흐름 이름들입니다. 단일 플로우 이름 또는 이름 리스트를 허용합니다. 이 명령은 오직 플로우 실행만 명령하며; 플로우 실행 방식은 변경하지 않습니다. depends_on와 함께 주문 파이프라인 흐름 실행을 참조하세요. |
comment |
str |
흐름에 대한 설명입니다. |
spark_conf |
dict |
이 쿼리 실행을 위한 Spark 구성 목록 |
import_checkpoint |
str |
기존 구조화 스트리밍 체크포인트로 가는 경로가 흐름에 가져오므로, 마이그레이션된 스트림이 소스를 재처리하지 않고 마지막으로 커밋된 오프셋에서 재개되도록 합니다. 체크포인트 가져오기는 베타 단계입니다. 구조 화된 스트리밍 체크포인트 마이그레이션을 참조하세요. |
예시
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"))
구조화된 스트리밍 체크포인트 마이그레이션
Important
체크포인트 가져오기는 베타 단계입니다.
기존 구조화 스트리밍 워크로드를 소스를 재처리하지 않고 파이프라인으로 마이그레이션하는 데 사용 import_checkpoint 하세요. 사용하는 구조화 스트리밍 쿼리(클라우드 저장소, Unity 카탈로그 볼륨, DBFS 경로)로 checkpointLocation 설정하세요. 첫 번째 파이프라인 업데이트에서는 해당 체크포인트를 파이프라인의 관리 저장소로 클론합니다. 흐름은 마지막으로 커밋된 오프셋에서 다시 시작되며, 그 상태(예: 집계, 중복 제거 키, 워터마크)가 그대로 유지됩니다. 이후 파이프라인 업데이트는 플로우의 복제된 체크포인트를 사용합니다; 원래 체크포인트는 수정되지 않습니다.
플로우는 create_table 또는 싱크로 생성된 관리된 테이블을 타겟팅해야 합니다.
파이프라인을 실행하기 전에 원래의 구조화된 스트리밍 쿼리를 멈추세요. 원래 구조화된 스트리밍 쿼리는 가져오기 후 재사용할 수 있지만, 체크포인트 상태를 관리하고 파이프라인과 쿼리가 동시에 같은 테이블에 쓰지 않도록 해야 합니다. 이는 중복 데이터가 발생할 수 있습니다.
구조화된 스트리밍 쿼리를 파이프라인 플로우로 재구성하여 새 테이블에 기록하고 체크포인트를 가져오세요:
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")
체크포인트는 첫 번째 파이프라인 업데이트에서 한 번만 가져오며; 이후 업데이트는 무시 import_checkpoint됩니다.
완전 새로고침은 체크포인트를 다시 가져오는 것이 아니며; 새롭고 빈 체크포인트에서 시작하여 소스를 재처리합니다. 다른 체크포인트를 가져오려면, 대상 테이블에 이전에 사용되지 않은 플로우 이름을 사용하세요; 기존 플로우 이름을 재사용하면 가져오기를 건너뛸 수 있습니다.
Limitations
- 이미 존재하는 테이블(예: 원래 Structured Streaming 쿼리 타겟)에 체크포인트를 가져오는 것은 지원되지 않습니다. 파이프라인이 생성하는 새 테이블이나 싱크를 타겟팅하세요.
-
import_checkpointappend_flow에서만 지원됩니다.