Pathway 핵심 기능 — 테이블·변환·데이터플로우
Pathway 핵심 기능
Pathway Live Data Framework는 파이프라인 정의와 실행을 명확히 분리해요. 파이프라인은 데이터 소스와 변환, 결과가 갈 곳을 정의하고, 실행은 그 파이프라인을 따라 실제 데이터를 처리해요. 핵심 개념을 하나씩 볼게요.
파이프라인 정의와 실행의 분리
파이프라인은 처리의 '레시피'예요. 재료(데이터 소스)와 여러 연산(변환)을 정의할 뿐, 실제 데이터(음식)를 담지는 않아요. 파이프라인을 구축한 뒤 실행해 데이터를 수집·처리해요.
입력 커넥터
커넥터는 외부 시스템과의 인터페이스예요. 입력 커넥터는 외부 데이터를 추출해 입력 스트림의 스냅샷을 나타내는 테이블을 반환해요. 스키마를 제공해야 해요.
import pathway as pw
class InputSchema(pw.Schema):
name: str = pw.column_definition(primary_key=True)
age: int
많은 입력 커넥터를 제공해요.
# CSV
input_table = pw.io.csv.read("./input_dir/", schema=InputSchema)
# Kafka
input_table = pw.io.kafka.read(rdkafka_settings, schema=InputSchema, topic="topic1", format="json")
# CDC (Debezium)
input_table = pw.io.debezium.read(rdkafka_settings, schema=InputSchema, topic_name="pets")
테이블: 정적 스키마, 동적 콘텐츠
데이터는 테이블로 모델링돼요. 관계형 테이블처럼 열·행으로 구성되고, 정적 스키마를 가지지만 콘텐츠는 동적이에요.
- 테이블은 데이터 스트림의 스냅샷, 현재 처리 시점까지 받은 이벤트의 최신 상태
- 각 새 항목은 한 행으로 (추가)
- 갱신·삭제도 지원 — 갱신은 이전 항목 삭제 + 새 버전 추가로 표현
name이 기본 키라서, 이미 있는name값을 가진 새 이벤트는 해당 행의 갱신이 돼요
변환 (Transformations)
select, join 같은 연산자로 테이블을 변환해요. 함수형 접근이라 각 변환은 새 테이블을 반환하고 입력 테이블은 그대로 두어요.
filtered_table = input_table.filter(input_table.age >= 0)
result_table = filtered_table.reduce(
sum_age = pw.reducers.sum(filtered_table.age)
)
주의: 이 코드는 아직 어떤 데이터 연산도 실행하지 않아요. 파이프라인을 정의만 하는 거예요.
외부 함수와 LLM
준비된 변환으로 부족하면 어떤 Python 함수든 파이프라인에 쓸 수 있어요. Python ML 라이브러리, LLM, 동기·비동기 API 모두 자연스럽게 통합돼요.
출력 커넥터
데이터는 출력 커넥터로 외부 시스템에 보내져요.
# CSV
pw.io.csv.write(table, "./output_file.csv")
# Kafka
pw.io.kafka.write(table, rdkafka_settings, topic="topic2", format="json")
# PostgreSQL
pw.io.postgres.write(table, psql_setting, "table_name")
증분 처리
Pathway는 데이터를 증분 방식으로 처리해요. 테이블 전체 버전을 보내기보다 변경된 행만 외부 시스템으로 보내 더 효율적이에요. 출력은 갱신 스트림으로 표현돼요(추가 diff=1, 삭제 diff=-1).