데코레이터는 @dp.foreach_batch_sink() 사용자 지정 논리를 사용하여 Python 처리하는 일련의 마이크로 일괄 처리로 스트림을 처리하는 ForEachBatch 싱크를 정의합니다. 싱크를 target의 싱크로 참조하여 변환된 데이터를 작성합니다. 개념 지침, 고려 사항 및 예제는 ForEachBatch를 사용하여 파이프라인의 임의 데이터 싱크에 쓰는 방법을 참조하세요.
Syntax
from pyspark import pipelines as dp
@dp.foreach_batch_sink(name="<name>")
def batch_handler(df, batch_id):
"""
Required:
- `df`: a Spark DataFrame representing the rows of this micro-batch.
- `batch_id`: unique integer ID for each micro-batch in the query.
"""
# Your custom write or transformation logic here
# Example:
# df.write.format("some-target-system").save("...")
#
# To access the sparkSession inside the batch handler, use df.sparkSession.
매개 변수
| 매개 변수 | Description |
|---|---|
| 이름 | Optional. 파이프라인 내에서 싱크를 식별하는 고유한 이름입니다. 포함되지 않은 경우 기본적으로 UDF의 이름으로 설정됩니다. |
| batch_handler | 각 마이크로 일괄 처리에 대해 호출되는 UDF(사용자 정의 함수)입니다. |
| df | 현재 마이크로 일괄 처리에 대한 데이터를 포함하는 Spark DataFrame입니다. |
| batch_id | 마이크로 배치의 정수 ID입니다. Spark는 각 트리거 간격에 대해 이 ID를 증가합니다.batch_id는 0 스트림의 시작이나 전체 새로 고침의 시작을 나타냅니다. 코드는 foreach_batch_sink 다운스트림 데이터 원본에 대한 전체 새로 고침을 제대로 처리해야 합니다. 자세한 내용은 전체 새로 고침을 참조하세요. |