DataStream API 튜토리얼

DataStream API 튜토리얼

Apache Flink 는 가장 낮은 수준의 스트림 처리 API 로 DataStream API 를 제공합니다. DataStream API 는 상태, 시간, 사용자 정의 처리 로직에 대한 세밀한 제어를 제공합니다. 이는 고급 이벤트 기반 애플리케이션을 구축하기에 이상적입니다.

출처: 문서

본문

Apache Flink 는 가장 낮은 수준의 스트림 처리 API 로 DataStream API 를 제공합니다. DataStream API 는 상태, 시간, 사용자 정의 처리 로직에 대한 세밀한 제어를 제공합니다. 이는 고급 이벤트 기반 애플리케이션을 구축하기에 이상적입니다.

무엇을 만들 것인가 (What You'll Build)

이 튜토리얼에서 여러분은 신용카드 거래를 처리하는 사기 탐지 시스템을 구축합니다:

Transactions (source) → Flink (KeyedProcessFunction) → Alerts (sink)

다음을 배우게 됩니다:

  • DataStream 실행 환경을 설정하는 방법
  • 데이터 수집을 위한 source 만드는 방법
  • 병렬 처리를 위해 keyBy 로 스트림을 파티셔닝하는 방법
  • KeyedProcessFunction 으로 비즈니스 로직을 구현하는 방법
  • 내결함성 있는 처리를 위해 관리되는 상태(ValueState)를 사용하는 방법

사전 요구사항 (Prerequisites)

이 워크스루는 Java 나 Python 에 대한 어느 정도의 친숙함을 가정하지만, 다른 프로그래밍 언어에서 왔더라도 따라올 수 있어야 합니다.

막혔을 때! (Help, I'm Stuck!)

막히면 커뮤니티 지원 리소스 를 확인하세요. 특히 Apache Flink 의 사용자 메일링 리스트 는 어떤 Apache 프로젝트보다도 활발한 곳으로 꾸준히 평가되며 빠르게 도움을 받는 좋은 방법입니다.

따라 하는 방법 (How To Follow Along)

따라 하려면 다음이 있는 컴퓨터가 필요합니다:

Java: Java 11, 17, 또는 21, Maven.

Python: Java 11, 17, 또는 21, Python 3.9, 3.10, 3.11, 또는 3.12.

Java: 제공된 Flink Maven Archetype 은 필요한 모든 의존성을 갖춘 스켈레톤 프로젝트를 빠르게 만들어 주므로 비즈니스 로직 작성에만 집중하면 됩니다. 이러한 의존성에는 모든 Flink 스트리밍 애플리케이션의 핵심 의존성인 flink-streaming-java 와 이 워크스루에 특화된 데이터 생성기 및 기타 클래스를 가진 flink-walkthrough-common 이 포함됩니다.

$ mvn archetype:generate \
    -DarchetypeGroupId=org.apache.flink \
    -DarchetypeArtifactId=flink-walkthrough-datastream-java \
    -DarchetypeVersion=2.3.0 \
    -DgroupId=frauddetection \
    -DartifactId=frauddetection \
    -Dversion=0.1 \
    -Dpackage=frauddetection \
    -DinteractiveMode=false

원한다면 groupId, artifactId, package 를 편집할 수 있습니다. 위 매개변수로 Maven 은 이 튜토리얼을 완료하는 데 필요한 모든 의존성을 가진 프로젝트가 들어 있는 frauddetection 이라는 폴더를 만듭니다.

프로젝트를 편집기에 가져온 후 IDE 에서 직접 실행할 수 있는 다음 코드가 있는 FraudDetectionJob.java 파일을 찾을 수 있습니다.

IDE 에서 실행: java.lang.NoClassDefFoundError 예외가 발생하면 클래스패스에 필요한 모든 Flink 의존성이 없기 때문일 가능성이 높습니다. IntelliJ IDEA 의 경우: Run > Edit Configurations > Modify options > "include dependencies with 'Provided' scope" 를 선택하세요.

Python: Python DataStream API 를 사용하려면 PyPI 에서 제공되고 pip 으로 쉽게 설치할 수 있는 PyFlink 를 설치해야 합니다:

$ python -m pip install apache-flink

팁: 프로젝트 의존성을 격리하기 위해 가상 환경(venv)에 PyFlink 를 설치할 것을 권장합니다.

PyFlink 가 설치되면 DataStream 프로그램을 작성할 fraud_detection.py 라는 새 파일을 만드세요.

전체 프로그램 (The Complete Program)

다음은 사기 탐지 프로그램의 전체 코드입니다:

Java (FraudDetectionJob):

public class FraudDetectionJob {

    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        DataStream<Transaction> transactions = env
            .fromSource(
                TransactionSource.unbounded(),
                WatermarkStrategy.noWatermarks(),
                "transactions");

        DataStream<Alert> alerts = transactions
            .keyBy(Transaction::getAccountId)
            .process(new FraudDetector())
            .name("fraud-detector");

        alerts
            .addSink(new AlertSink())
            .name("send-alerts");

        env.execute("Fraud Detection");
    }
}

Java (FraudDetector - 첫 버전):

public class FraudDetector extends KeyedProcessFunction<Long, Transaction, Alert> {

    private static final long serialVersionUID = 1L;

    private static final double SMALL_AMOUNT = 1.00;
    private static final double LARGE_AMOUNT = 500.00;

    private transient ValueState<Boolean> flagState;

    @Override
    public void open(OpenContext openContext) {
        ValueStateDescriptor<Boolean> flagDescriptor = new ValueStateDescriptor<>(
                "flag",
                Types.BOOLEAN);
        flagState = getRuntimeContext().getState(flagDescriptor);
    }

    @Override
    public void processElement(
            Transaction transaction,
            Context context,
            Collector<Alert> collector) throws Exception {

        Boolean lastTransactionWasSmall = flagState.value();

        if (lastTransactionWasSmall != null) {
            if (transaction.getAmount() > LARGE_AMOUNT) {
                Alert alert = new Alert();
                alert.setId(transaction.getAccountId());
                collector.collect(alert);
            }
            flagState.clear();
        }

        if (transaction.getAmount() < SMALL_AMOUNT) {
            flagState.update(true);
        }
    }
}

Python (첫 버전):

from pyflink.common.typeinfo import Types
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.datastream.functions import KeyedProcessFunction, RuntimeContext
from pyflink.datastream.state import ValueStateDescriptor


class FraudDetector(KeyedProcessFunction):

    SMALL_AMOUNT = 1.00
    LARGE_AMOUNT = 500.00

    def __init__(self):
        self.flag_state = None

    def open(self, runtime_context: RuntimeContext):
        descriptor = ValueStateDescriptor("flag", Types.BOOLEAN())
        self.flag_state = runtime_context.get_state(descriptor)

    def process_element(self, transaction, ctx: 'KeyedProcessFunction.Context'):
        # transaction is a tuple: (account_id, timestamp, amount)
        account_id = transaction[0]
        amount = transaction[2]

        last_transaction_was_small = self.flag_state.value()

        if last_transaction_was_small is not None:
            if amount > self.LARGE_AMOUNT:
                yield f"Alert{{id={account_id}}}"
            self.flag_state.clear()

        if amount < self.SMALL_AMOUNT:
            self.flag_state.update(True)


def fraud_detection():
    env = StreamExecutionEnvironment.get_execution_environment()
    env.set_parallelism(1)

    # Sample transaction data: (account_id, timestamp, amount)
    transactions_data = [
        (1, 1000, 188.23),
        (2, 1001, 0.50),    # Small transaction
        (2, 1002, 600.00),  # Large transaction - ALERT!
        (3, 1003, 42.00),
        (1, 1004, 0.89),    # Small transaction
        (1, 1005, 300.00),  # Not large enough - no alert
        (4, 1006, 0.10),    # Small transaction
        (4, 1007, 520.00),  # Large transaction - ALERT!
        (3, 1008, 0.75),    # Small transaction
        (3, 1009, 800.00),  # Large transaction - ALERT!
    ]

    transactions = env.from_collection(
        transactions_data,
        type_info=Types.TUPLE([Types.LONG(), Types.LONG(), Types.DOUBLE()])
    )

    alerts = transactions \
        .key_by(lambda t: t[0]) \
        .process(FraudDetector())

    alerts.print()

    env.execute("Fraud Detection")


if __name__ == '__main__':
    fraud_detection()

코드 분석 (Breaking Down The Code)

코드를 단계별로 살펴보겠습니다. main 클래스는 애플리케이션의 데이터 흐름을 정의하고 FraudDetector 클래스는 사기 거래를 감지하는 함수의 비즈니스 로직을 정의합니다.

실행 환경 (The Execution Environment)

첫 번째 줄은 StreamExecutionEnvironment 를 설정합니다. 실행 환경은 Job 의 속성을 설정하고, source 를 만들고, Job 의 실행을 트리거하는 방법입니다.

Java:

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

Python:

env = StreamExecutionEnvironment.get_execution_environment()

Source 만들기

Source 는 Apache Kafka, Rabbit MQ, Apache Pulsar 같은 외부 시스템의 데이터를 Flink Job 으로 수집합니다.

Java: 이 워크스루는 처리할 신용카드 거래의 무한 스트림을 생성하는 DataGeneratorSource 를 감싼 TransactionSource 를 사용합니다. 각 거래에는 계정 ID(accountId), 거래 발생 시각(timestamp), 금액(amount) 이 포함됩니다.

DataStream<Transaction> transactions = env
    .fromSource(
        TransactionSource.unbounded(),
        WatermarkStrategy.noWatermarks(),
        "transactions");

fromSource 메서드는 세 가지 매개변수를 받습니다: source 자체, 워터마크 전략(이 예제는 processing time 을 사용하므로 noWatermarks() 사용), 디버깅용 이름입니다.

Python: 이 워크스루는 샘플 거래 데이터 모음을 사용합니다. 각 거래는 계정 ID, 타임스탬프, 금액을 포함하는 튜플입니다.

transactions_data = [
    (1, 1000, 188.23),
    (2, 1001, 0.50),    # Small transaction
    (2, 1002, 600.00),  # Large transaction - ALERT!
    # ...more transactions
]

transactions = env.from_collection(
    transactions_data,
    type_info=Types.TUPLE([Types.LONG(), Types.LONG(), Types.DOUBLE()])
)

프로덕션 시스템에서는 Kafka 같은 소스 커넥터를 사용할 것입니다. from_collection 메서드는 테스트와 튜토리얼에 편리합니다.

이벤트 파티셔닝 및 사기 탐지

transactions 스트림은 많은 사용자의 수많은 거래를 포함하므로 여러 사기 탐지 작업이 병렬로 처리해야 합니다. 사기는 계정 단위로 발생하므로, 같은 계정의 모든 거래가 사기 탐지 operator 의 같은 병렬 작업에 의해 처리되도록 보장해야 합니다.

같은 물리적 작업이 특정 키의 모든 레코드를 처리하도록 보장하려면 keyBy 를 사용해 스트림을 파티셔닝할 수 있습니다. process() 호출은 스트림의 각 파티셔닝된 요소에 함수를 적용하는 operator 를 추가합니다. keyBy 바로 다음의 operator, 여기서는 FraudDetectorkeyed context 내에서 실행된다고 말하는 것이 일반적입니다.

Java:

DataStream<Alert> alerts = transactions
    .keyBy(Transaction::getAccountId)
    .process(new FraudDetector())
    .name("fraud-detector");

Python:

alerts = transactions \
    .key_by(lambda t: t[0]) \
    .process(FraudDetector())

결과 출력

sink 는 DataStream 을 Apache Kafka, Cassandra, AWS Kinesis 같은 외부 시스템에 씁니다.

Java: AlertSink 는 영구 저장소에 쓰는 대신 각 Alert 레코드를 INFO 로그 레벨로 기록하므로 결과를 쉽게 볼 수 있습니다.

alerts.addSink(new AlertSink());

Python: print() 메서드는 쉽게 볼 수 있도록 경보를 콘솔에 출력합니다.

alerts.print()

사기 탐지기 (The Fraud Detector)

사기 탐지기는 KeyedProcessFunction 으로 구현됩니다. 그 processElement 메서드는 모든 거래 이벤트에 대해 호출됩니다. 이 첫 버전은 모든 거래에 대해 경보를 생성하며, 과도하게 보수적이라고 볼 수 있습니다. 다음 섹션은 더 의미 있는 비즈니스 로직으로 사기 탐지기를 확장하도록 안내합니다.

Java:

public class FraudDetector extends KeyedProcessFunction<Long, Transaction, Alert> {

    private static final double SMALL_AMOUNT = 1.00;
    private static final double LARGE_AMOUNT = 500.00;

    @Override
    public void processElement(
            Transaction transaction,
            Context context,
            Collector<Alert> collector) throws Exception {

        Alert alert = new Alert();
        alert.setId(transaction.getAccountId());

        collector.collect(alert);
    }
}

Python:

class FraudDetector(KeyedProcessFunction):

    SMALL_AMOUNT = 1.00
    LARGE_AMOUNT = 500.00

    def process_element(self, transaction, ctx: 'KeyedProcessFunction.Context'):
        account_id = transaction[0]
        yield f"Alert{{id={account_id}}}"

비즈니스 로직 구현하기

첫 버전의 경우 사기 탐지기는 작은 거래 직후에 큰 거래가 이어지는 모든 계정에 경보를 출력해야 합니다. 여기서 small 은 1.00 미만, large 는 500 초과입니다.

사기 탐지기가 특정 계정에 대해 다음 거래 스트림을 처리한다고 상상해 보세요.

거래 3과 4는 작은 거래(0.09) 직후에 큰 거래(510)가 이어지므로 사기로 표시되어야 합니다. 반면 거래 7, 8, 9는 작은 금액(0.02) 직후에 큰 거래가 바로 이어지지 않았으므로(중간에 패턴을 깨는 거래가 있음) 사기가 아닙니다.

이를 수행하려면 사기 탐지기가 이벤트 간 정보를 기억 해야 합니다; 큰 거래는 이전 거래가 작았을 때만 사기입니다. 이벤트 간 정보를 기억하려면 상태 가 필요하며, 그래서 이 튜토리얼은 KeyedProcessFunction 을 사용합니다. 이것은 상태와 시간 모두에 대한 세밀한 제어를 제공하여 더 복잡한 요구사항으로 알고리즘을 발전시킬 수 있게 합니다.

가장 간단한 구현은 작은 거래가 처리될 때마다 설정되는 boolean 플래그입니다. 큰 거래가 들어오면 해당 계정에 플래그가 설정되었는지 간단히 확인할 수 있습니다.

그러나 플래그를 단순히 FraudDetector 클래스의 멤버 변수로 구현하는 것은 작동하지 않습니다. Flink 는 FraudDetector 의 같은 객체 인스턴스로 여러 계정의 거래를 처리합니다. 즉 계정 A 와 B 가 같은 FraudDetector 인스턴스로 라우팅되면 계정 A 의 거래가 플래그를 true 로 설정한 후 계정 B 의 거래가 잘못된 경보를 유발할 수 있습니다. 물론 Map 같은 데이터 구조를 사용해 개별 키의 플래그를 추적할 수 있지만, 단순 멤버 변수는 내결함성이 없어 실패 시 모든 정보가 손실됩니다. 따라서 애플리케이션이 실패에서 복구하기 위해 재시작해야 한다면 사기 탐지기가 경보를 놓칠 수 있습니다.

이러한 문제를 해결하기 위해 Flink 는 일반 멤버 변수만큼 사용하기 쉬운 내결함성 상태를 위한 프리미티브를 제공합니다.

Flink 의 가장 기본적인 상태 유형은 ValueState 로, 감싸는 어떤 변수에도 내결함성을 추가하는 데이터 타입입니다. ValueStatekeyed state 의 한 형태로, keyed context 에 적용된 operator 에서만 사용할 수 있습니다; 즉 keyBy 바로 다음의 operator 입니다. operator 의 keyed state 는 현재 처리 중인 레코드의 키로 자동 범위가 지정됩니다. 이 예제에서 키는 현재 거래의 계정 id(keyBy() 로 선언됨) 이며, FraudDetector 는 각 계정에 대해 독립적인 상태를 유지합니다.

ValueState 는 Flink 가 변수를 어떻게 관리해야 하는지에 대한 메타데이터를 포함하는 ValueStateDescriptor 를 사용해 생성됩니다. 상태는 함수가 데이터 처리를 시작하기 전에 등록되어야 합니다. 그에 맞는 훅은 open() 메서드입니다.

Java:

public class FraudDetector extends KeyedProcessFunction<Long, Transaction, Alert> {

    private static final long serialVersionUID = 1L;

    private transient ValueState<Boolean> flagState;

    @Override
    public void open(OpenContext openContext) {
        ValueStateDescriptor<Boolean> flagDescriptor = new ValueStateDescriptor<>(
                "flag",
                Types.BOOLEAN);
        flagState = getRuntimeContext().getState(flagDescriptor);
    }

Python:

class FraudDetector(KeyedProcessFunction):

    def __init__(self):
        self.flag_state = None

    def open(self, runtime_context: RuntimeContext):
        descriptor = ValueStateDescriptor("flag", Types.BOOLEAN())
        self.flag_state = runtime_context.get_state(descriptor)

ValueState 는 Java 표준 라이브러리의 AtomicReferenceAtomicLong 과 유사한 래퍼 클래스입니다. 내용과 상호작용하는 세 가지 메서드를 제공합니다: update 는 상태를 설정하고, value 는 현재 값을 얻고, clear 는 내용을 삭제합니다. 특정 키의 상태가 비어 있으면(애플리케이션 시작 시나 clear 호출 후처럼) valuenull (Java) 또는 None (Python) 을 반환합니다. value 가 반환한 객체에 대한 수정은 시스템이 인식한다는 보장이 없으므로 모든 변경은 update 로 수행해야 합니다. 그 외에는 내결함성이 Flink 에 의해 내부적으로 자동 관리되므로 다른 표준 변수처럼 상호작용할 수 있습니다.

아래는 플래그 상태를 사용해 잠재적 사기 거래를 추적하는 예입니다:

Java:

@Override
public void processElement(
        Transaction transaction,
        Context context,
        Collector<Alert> collector) throws Exception {

    // Get the current state for the current key
    Boolean lastTransactionWasSmall = flagState.value();

    // Check if the flag is set
    if (lastTransactionWasSmall != null) {
        if (transaction.getAmount() > LARGE_AMOUNT) {
            // Output an alert downstream
            Alert alert = new Alert();
            alert.setId(transaction.getAccountId());

            collector.collect(alert);
        }
        // Clean up our state
        flagState.clear();
    }

    if (transaction.getAmount() < SMALL_AMOUNT) {
        // Set the flag to true
        flagState.update(true);
    }
}

Python:

def process_element(self, transaction, ctx: 'KeyedProcessFunction.Context'):
    account_id = transaction[0]
    amount = transaction[2]

    # Get the current state for the current key
    last_transaction_was_small = self.flag_state.value()

    # Check if the flag is set
    if last_transaction_was_small is not None:
        if amount > self.LARGE_AMOUNT:
            # Output an alert downstream
            yield f"Alert{{id={account_id}}}"

        # Clean up our state
        self.flag_state.clear()

    if amount < self.SMALL_AMOUNT:
        # Set the flag to true
        self.flag_state.update(True)

모든 거래에 대해 사기 탐지기는 해당 계정의 플래그 상태를 확인합니다. ValueState 는 항상 현재 키, 즉 계정으로 범위가 지정된다는 것을 기억하세요. 플래그가 null 이 아니면 해당 계정에 대해 본 마지막 거래가 작은 것이므로, 이 거래의 금액이 크면 탐지기가 사기 경보를 출력합니다.

그 확인 후 플래그 상태는 무조건 지워집니다. 현재 거래가 사기 경보를 일으켰다면 패턴이 끝난 것이고, 현재 거래가 경보를 일으키지 않았다면 패턴이 깨졌으므로 다시 시작해야 합니다.

마지막으로 거래 금액이 작은지 확인합니다. 작으면 플래그가 설정되어 다음 이벤트가 확인할 수 있게 합니다. ValueState<Boolean> 은 모든 ValueState 가 nullable 이므로 unset(null/None), true/True, false/False 의 세 가지 상태를 가집니다. 이 작업은 플래그가 설정되었는지 확인하기 위해 unset 과 true 만 사용합니다.

완전한 구현 (Complete Implementation)

다음은 상태 기반 사기 탐지가 포함된 완전한 FraudDetector 구현입니다:

Java:

import org.apache.flink.api.common.functions.OpenContext;
import org.apache.flink.api.common.state.ValueState;
import org.apache.flink.api.common.state.ValueStateDescriptor;
import org.apache.flink.api.common.typeinfo.Types;
import org.apache.flink.streaming.api.functions.KeyedProcessFunction;
import org.apache.flink.util.Collector;
import org.apache.flink.walkthrough.common.entity.Alert;
import org.apache.flink.walkthrough.common.entity.Transaction;

public class FraudDetector extends KeyedProcessFunction<Long, Transaction, Alert> {

    private static final long serialVersionUID = 1L;

    private static final double SMALL_AMOUNT = 1.00;
    private static final double LARGE_AMOUNT = 500.00;

    private transient ValueState<Boolean> flagState;

    @Override
    public void open(OpenContext openContext) {
        ValueStateDescriptor<Boolean> flagDescriptor = new ValueStateDescriptor<>(
                "flag",
                Types.BOOLEAN);
        flagState = getRuntimeContext().getState(flagDescriptor);
    }

    @Override
    public void processElement(
            Transaction transaction,
            Context context,
            Collector<Alert> collector) throws Exception {

        // Get the current state for the current key
        Boolean lastTransactionWasSmall = flagState.value();

        // Check if the flag is set
        if (lastTransactionWasSmall != null) {
            if (transaction.getAmount() > LARGE_AMOUNT) {
                // Output an alert downstream
                Alert alert = new Alert();
                alert.setId(transaction.getAccountId());

                collector.collect(alert);
            }
            // Clean up our state
            flagState.clear();
        }

        if (transaction.getAmount() < SMALL_AMOUNT) {
            // Set the flag to true
            flagState.update(true);
        }
    }
}

Python:

from pyflink.common.typeinfo import Types
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.datastream.functions import KeyedProcessFunction, RuntimeContext
from pyflink.datastream.state import ValueStateDescriptor


class FraudDetector(KeyedProcessFunction):

    SMALL_AMOUNT = 1.00
    LARGE_AMOUNT = 500.00

    def __init__(self):
        self.flag_state = None

    def open(self, runtime_context: RuntimeContext):
        descriptor = ValueStateDescriptor("flag", Types.BOOLEAN())
        self.flag_state = runtime_context.get_state(descriptor)

    def process_element(self, transaction, ctx: 'KeyedProcessFunction.Context'):
        # transaction is a tuple: (account_id, timestamp, amount)
        account_id = transaction[0]
        amount = transaction[2]

        # Get the current state for the current key
        last_transaction_was_small = self.flag_state.value()

        # Check if the flag is set
        if last_transaction_was_small is not None:
            if amount > self.LARGE_AMOUNT:
                # Output an alert downstream
                yield f"Alert{{id={account_id}}}"

            # Clean up our state
            self.flag_state.clear()

        if amount < self.SMALL_AMOUNT:
            # Set the flag to true
            self.flag_state.update(True)


def fraud_detection():
    env = StreamExecutionEnvironment.get_execution_environment()
    env.set_parallelism(1)

    # Sample transaction data: (account_id, timestamp, amount)
    transactions_data = [
        (1, 1000, 188.23),
        (2, 1001, 0.50),    # Small transaction
        (2, 1002, 600.00),  # Large transaction - ALERT!
        (3, 1003, 42.00),
        (1, 1004, 0.89),    # Small transaction
        (1, 1005, 300.00),  # Not large enough - no alert
        (4, 1006, 0.10),    # Small transaction
        (4, 1007, 520.00),  # Large transaction - ALERT!
        (3, 1008, 0.75),    # Small transaction
        (3, 1009, 800.00),  # Large transaction - ALERT!
    ]

    transactions = env.from_collection(
        transactions_data,
        type_info=Types.TUPLE([Types.LONG(), Types.LONG(), Types.DOUBLE()])
    )

    alerts = transactions \
        .key_by(lambda t: t[0]) \
        .process(FraudDetector())

    alerts.print()

    env.execute("Fraud Detection")


if __name__ == '__main__':
    fraud_detection()

애플리케이션 실행

이제 끝입니다. 완전히 기능하는, 상태 기반의 분산 스트리밍 애플리케이션이 완성되었습니다! 쿼리는 소스에서 거래를 지속적으로 소비하고, 사기 패턴을 감지하며, 경보를 내보냅니다. (Java 버전에서는) 입력이 무한하므로 쿼리는 수동으로 중지될 때까지 계속 실행됩니다.

Java: IDE 에서 FraudDetectionJob 클래스를 실행해 스트리밍 결과를 확인하세요. 다음과 유사한 출력이 보여야 합니다:

2024-01-01 14:22:06,220 INFO  org.apache.flink.walkthrough.common.sink.AlertSink - Alert{id=3}
2024-01-01 14:22:11,383 INFO  org.apache.flink.walkthrough.common.sink.AlertSink - Alert{id=3}
2024-01-01 14:22:16,551 INFO  org.apache.flink.walkthrough.common.sink.AlertSink - Alert{id=3}

Python: 명령줄에서 프로그램을 실행하세요:

$ python fraud_detection.py

다음과 유사한 출력이 보여야 합니다:

Alert{id=2}
Alert{id=4}
Alert{id=3}

이 명령은 로컬 미니 클러스터에서 PyFlink 프로그램을 빌드하고 실행합니다. Python DataStream 프로그램을 원격 클러스터에 제출할 수도 있습니다. 자세한 내용은 Job Submission Examples 를 참고하세요.

다음 단계 (Next Steps)

튜토리얼을 완료한 것을 축하합니다! 학습을 계속할 몇 가지 방법은 다음과 같습니다:

사기 탐지기 향상

현재 구현은 작은-후-큰 거래 패턴을 감지하지만, 실제 사기 탐지에는 시간 제약이 포함되는 경우가 많습니다. 예를 들어 사기범은 보통 테스트 거래와 큰 구매 사이를 오래 기다리지 않습니다.

사기 탐지기에 타이머를 추가하는 방법(예: 서로 1분 이내에 발생한 거래만 표시)은 Learn Flink 의 Event-Driven Applications 섹션을 참고하세요.

DataStream 에 대해 더 배우기

다른 튜토리얼 탐색

프로덕션 배포

더 알아보기 (Learn more)