append_flow

@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 다음 두 가지 방법으로 흐름을 변경합니다.
  • 반환 값입니다. streaming-query; 이 경우는 스트리밍 데이터 프레임이 아닌 배치 DataFrame이어야 합니다.
  • 흐름은 기본적으로 한 번 실행됩니다. 파이프라인이 전체 새로 고침으로 업데이트되면 흐름이 ONCE 다시 실행되어 데이터를 다시 만듭니다.
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_checkpoint append_flow에서만 지원됩니다.