외부 데이터 접근을 위한 비동기 I/O

외부 데이터 접근을 위한 비동기 I/O (Asynchronous I/O for External Data Access)

이 페이지는 외부 데이터 저장소와의 비동기 I/O를 위한 Flink API 사용법을 설명해요. 비동기 또는 이벤트 주도 프로그래밍에 익숙하지 않은 사용자에게는 Futures와 이벤트 주도 프로그래밍에 관한 글이 유용한 준비가 될 수 있어요.

출처: 문서

본문

참고: 비동기 I/O 유틸리티의 설계·구현 세부사항은 제안 및 설계 문서 FLIP-12: Asynchronous I/O Design and Implementation에서 찾을 수 있어요. 새 재시도 지원 세부사항은 FLIP-232: Add Retry Support For Async I/O In DataStream API 문서에서 찾을 수 있어요.

비동기 I/O 연산의 필요성 (The need for Asynchronous I/O Operations)

외부 시스템과 상호작용할 때(예: 데이터베이스에 저장된 데이터로 스트림 이벤트를 보강할 때) 외부 시스템과의 통신 지연이 스트리밍 애플리케이션의 전체 작업을 지배하지 않도록 주의해야 해요. 외부 데이터베이스의 데이터에 단순히 접근하는 것(예: MapFunction에서)은 보통 동기적 상호작용을 의미해요. 요청을 데이터베이스로 보내고 MapFunction이 응답을 받을 때까지 기다려요. 많은 경우 이 대기는 함수 시간의 대부분을 차지해요.

데이터베이스와의 비동기 상호작용은 단일 병렬 함수 인스턴스가 많은 요청을 동시에 처리하고 응답도 동시에 받을 수 있음을 의미해요. 그렇게 하면 대기 시간을 다른 요청 보내기와 응답 받기에 겹칠 수 있어요. 적어도 대기 시간은 여러 요청에 걸쳐 상각(amortize)되어요. 이는 대부분의 경우 훨씬 더 높은 스트리밍 처리량으로 이어져요.

참고: MapFunction을 매우 높은 병렬도로 확장해 처리량을 개선하는 것도 일부 경우 가능하지만, 보통 매우 높은 리소스 비용이 들어요. 병렬 MapFunction 인스턴스가 훨씬 많다는 것은 더 많은 작업, 스레드, Flink 내부 네트워크 연결, 데이터베이스로의 네트워크 연결, 버퍼, 일반적인 내부 부기 오버헤드를 의미해요.

전제 조건 (Prerequisites)

위 섹션에서 설명했듯이 데이터베이스(또는 키/값 저장소)에 대한 적절한 비동기 I/O를 구현하려면 비동기 요청을 지원하는 데이터베이스 클라이언트가 필요해요. 많은 인기 데이터베이스가 그러한 클라이언트를 제공해요. 그런 클라이언트가 없다면 여러 클라이언트를 만들고 스레드 풀로 동기 호출을 처리해 동기 클라이언트를 제한된 동시 클라이언트로 바꿔볼 수 있어요. 그러나 이 접근은 보통 적절한 비동기 클라이언트보다 덜 효율적이에요.

Async I/O API

Flink의 Async I/O API는 사용자가 데이터 스트림과 함께 비동기 요청 클라이언트를 사용할 수 있게 해줘요. 이 API는 데이터 스트림과의 통합뿐 아니라 순서, 이벤트 시간, 장애 허용, 재시도 지원 등을 처리해요.

대상 데이터베이스용 비동기 클라이언트가 있다고 가정하면, 데이터베이스에 대한 비동기 I/O 스트림 변환을 구현하려면 세 부분이 필요해요:

  • 요청을 파견(dispatch)하는 AsyncFunction 구현
  • 연산의 결과를 받아 Java API에서는 ResultFuture에 넘기거나 Python API에서는 연산 결과를 기다리는 콜백
  • 재시도 유무와 관계없이 DataStream에 async I/O 연산을 변환으로 적용

다음 코드 예제는 기본 패턴을 보여줘요.

Java:

// This example implements the asynchronous request and callback with Futures that have the
// interface of Java 8's futures (which is the same one followed by Flink's Future)

/**
* An implementation of the 'AsyncFunction' that sends requests and sets the callback.
*/
class AsyncDatabaseRequest extends RichAsyncFunction<String, Tuple2<String, String>> {

/** The database specific client that can issue concurrent requests with callbacks */
private transient DatabaseClient client;

@Override
public void open(OpenContext openContext) throws Exception {
client = new DatabaseClient(host, post, credentials);
}

@Override
public void close() throws Exception {
client.close();
}

@Override
public void asyncInvoke(String key, final ResultFuture<Tuple2<String, String>> resultFuture) throws Exception {

// issue the asynchronous request, receive a future for result
final Future<String> result = client.query(key);

// set the callback to be executed once the request by the client is complete
// the callback simply forwards the result to the result future
CompletableFuture.supplyAsync(new Supplier<String>() {

@Override
public String get() {
try {
return result.get();
} catch (InterruptedException | ExecutionException e) {
// Normally handled explicitly.
return null;
}
}
}).thenAccept( (String dbResult) -> {
resultFuture.complete(Collections.singleton(new Tuple2<>(key, dbResult)));
});
}
}

// create the original stream
DataStream<String> stream = ...;

// apply the async I/O transformation without retry
DataStream<Tuple2<String, String>> resultStream =
AsyncDataStream.unorderedWait(stream, new AsyncDatabaseRequest(), 1000, TimeUnit.MILLISECONDS, 100);

// or apply the async I/O transformation with retry
// create an async retry strategy via utility class or a user defined strategy
AsyncRetryStrategy asyncRetryStrategy =
new AsyncRetryStrategies.FixedDelayRetryStrategyBuilder(3, 100L) // maxAttempts=3, fixedDelay=100ms
.ifResult(RetryPredicates.EMPTY_RESULT_PREDICATE)
.ifException(RetryPredicates.HAS_EXCEPTION_PREDICATE)
.build();

// apply the async I/O transformation with retry
DataStream<Tuple2<String, String>> resultStream =
AsyncDataStream.unorderedWaitWithRetry(stream, new AsyncDatabaseRequest(), 1000, TimeUnit.MILLISECONDS, 100, asyncRetryStrategy);

Python:

from typing import List

from pyflink.common import Time, Types
from pyflink.datastream import AsyncFunction, AsyncDataStream, async_retry_predicates
from pyflink.datastream.functions import RuntimeContext, AsyncRetryStrategy

class AsyncDatabaseRequest(AsyncFunction[str, (str, str)]):

def __init__(self, host, port, credentials):
self._host = host
self._port = port
self._credentials = credentials

def open(self, runtime_context: RuntimeContext):
# The database specific client that can issue concurrent requests with callbacks
self._client = DatabaseClient(self._host, self._port, self._credentials)

def close(self):
if self._client:
self._client.close()

async def async_invoke(self, value: str) -> List[(str, str)]:
try:
# issue the asynchronous request
result = await self._client.query(value)
return [(value, str(result))]
except Exception:
return [(value, None)]

# create the original stream
stream = ...

# apply the async I/O transformation without retry
result_stream = AsyncDataStream.unordered_wait(
data_stream=stream,
async_function=AsyncDatabaseRequest("127.0.0.1", "1234", None),
timeout=Time.seconds(10),
capacity=100,
output_type=Types.TUPLE([Types.STRING(), Types.STRING()]))

# or apply the async I/O transformation with retry
# create an async retry strategy via utility class or a user defined strategy
async_retry_strategy = AsyncRetryStrategy.fixed_delay(
max_attempts=3,
backoff_time_millis=100,
result_predicate=async_retry_predicates.empty_result_predicate,
exception_predicate=async_retry_predicates.has_exception_predicate)

# apply the async I/O transformation with retry
result_stream_with_retry = AsyncDataStream.unordered_wait_with_retry(
data_stream=stream,
async_function=AsyncDatabaseRequest("127.0.0.1", "1234", None),
timeout=Time.seconds(10),
async_retry_strategy=async_retry_strategy,
capacity=1000,
output_type=Types.TUPLE([Types.STRING(), Types.STRING()]))

중요 참고: Java API에서 ResultFuture는 ResultFuture.complete의 첫 호출로 완료돼요. 이후의 모든 complete 호출은 무시돼요.

다음 세 파라미터가 비동기 연산을 제어해요:

  • Timeout (타임아웃): 타임아웃은 비동기 연산의 첫 호출부터 최종 완료까지의 최대 기간을 정의해요. 이 기간은 여러 재시도 시도(재시도가 활성화된 경우)를 포함할 수 있으며 연산이 궁극적으로 완료된 것으로 간주되는 시점을 결정해요. 이 파라미터는 죽었거나 실패한 요청을 보호해요.
  • Capacity (용량): 이 파라미터는 async 연산자의 병렬 인스턴스(하위 작업)당 동시에 진행 중일 수 있는 비동기 요청 수를 정의해요. 비동기 I/O 접근이 일반적으로 훨씬 더 나은 처리량을 제공하지만 연산자는 여전히 스트리밍 애플리케이션의 병목일 수 있어요. 동시 요청 수를 제한하면 연산자가 계속 커지는 보류 요청 백로그를 쌓지 않고, 용량이 소진되면 백프레셔를 트리거함을 보장해요.
  • AsyncRetryStrategy: 이 파라미터는 어떤 조건이 지연된 재시도를 트리거할지와 지연 전략(예: 고정 지연, 지수 백오프 지연, 사용자 정의 구현 등)을 정의해요.

타임아웃 처리 (Timeout Handling)

비동기 I/O 요청이 타임아웃되면 기본적으로 예외가 발생하고 작업이 재시작돼요. 타임아웃을 처리하려면 AsyncFunction#timeout 메서드를 오버라이드할 수 있어요. Java API에서는 오버라이드할 때 이 입력 레코드의 처리가 완료되었음을 Flink에 알리기 위해 ResultFuture.complete() 또는 ResultFuture.completeExceptionally()를 호출해야 해요. 타임아웃 시 레코드를 발행하지 않으려면 ResultFuture.complete(Collections.emptyList())를 호출할 수 있어요. Python API에서는 오버라이드할 때 결과 컬렉션을 반환하거나 예외를 발생시켜 이 입력 레코드의 처리가 완료되었음을 Flink에 알릴 수 있어요. 타임아웃 시 레코드를 발행하지 않으려면 return []로 빈 목록을 반환할 수 있어요.

결과 순서 (Order of Results)

AsyncFunction이 발행한 동시 요청은 어떤 요청이 먼저 끝났는지에 따라 종종 정의되지 않은 순서로 완료돼요. 결과 레코드가 발행되는 순서를 제어하기 위해 Flink는 두 가지 모드를 제공해요:

  • Unordered: 결과 레코드는 비동기 요청이 끝나는 즉시 발행돼요. async I/O 연산자 이후 스트림의 레코드 순서는 이전과 달라져요. 이 모드는 processing time을 기본 시간 특성으로 사용할 때 가장 낮은 지연과 오버헤드를 가져요. 이 모드에는 AsyncDataStream.unorderedWait(...) 또는 AsyncDataStream.unordered_wait(...)을 사용해요.
  • Ordered: 이 경우 스트림 순서가 보존돼요. 결과 레코드는 비동기 요청이 트리거된 순서(연산자 입력 레코드의 순서)와 같은 순서로 발행돼요. 이를 위해 연산자는 앞선 모든 레코드가 발행(또는 타임아웃)될 때까지 결과 레코드를 버퍼링해요. 이는 보통 unordered 모드와 비교해 레코드나 결과가 더 오래 체크포인트 상태에 유지되므로 약간의 추가 지연과 체크포인팅 오버헤드를 도입해요. 이 모드에는 AsyncDataStream.orderedWait(...) 또는 AsyncDataStream.ordered_wait(...)을 사용해요.

이벤트 시간 (Event Time)

스트리밍 애플리케이션이 event time으로 동작할 때 워터마크는 비동기 I/O 연산자에 의해 올바르게 처리돼요. 이는 두 순서 모드에 대해 구체적으로 다음을 의미해요:

  • Unordered: 워터마크는 레코드를 추월하지 않고 그 반대도 마찬가지이며, 워터마크가 순서 경계를 만든다는 뜻이에요. 레코드는 워터마크 사이에서만 순서 없이 발행돼요. 특정 워터마크 이후에 발생한 레코드는 그 워터마크가 발행된 후에만 발행돼요. 워터마크는 그 워터마크 이전의 입력의 모든 결과 레코드가 발행된 후에만 발행돼요. 즉 워터마크가 있으면 unordered 모드도 ordered 모드와 같은 지연과 관리 오버헤드를 일부 도입해요. 그 오버헤드의 양은 워터마크 빈도에 따라 달라져요.
  • Ordered: 레코드 간 순서가 보존되는 것처럼 워터마크와 레코드의 순서가 보존돼요. processing time으로 작업하는 것과 비교해 오버헤드의 큰 변화는 없어요.

Ingestion Time은 소스의 처리 시간에 기반한 자동 생성 워터마크를 가진 event time의 특수한 경우임을 기억해요.

장애 허용 보장 (Fault Tolerance Guarantees)

비동기 I/O 연산자는 완전한 exactly-once 장애 허용 보장을 제공해요. 진행 중인 비동기 요청의 레코드를 체크포인트에 저장하고 실패에서 복구할 때 요청을 복원/재트리거해요.

재시도 지원 (Retry Support)

재시도 지원은 사용자의 AsyncFunction에 투명한 async 연산자용 내장 메커니즘을 도입해요.

  • AsyncRetryStrategy: AsyncRetryStrategy는 재시도 조건 AsyncRetryPredicate의 정의와 현재 시도 횟수에 따라 재시도를 계속할지와 재시도 간격을 결정하는 인터페이스를 포함해요. 재시도 트리거 조건이 충족된 후에도 현재 시도 횟수가 사전 설정 한도를 초과하면 재시도를 포기할 수 있고, 또는 작업이 끝날 때 재시도를 강제로 종료할 수 있어요(이 경우 시스템은 마지막 실행 결과 또는 예외를 최종 상태로 삼아요).
  • AsyncRetryPredicate: 재시도 조건은 반환 결과 또는 실행 예외에 기반해 트리거될 수 있어요.

구현 팁 (Implementation Tips)

콜백용 Executor가 있는 Futures 구현의 경우 DirectExecutor를 사용할 것을 권장해요. 콜백은 보통 최소한의 작업을 하므로 DirectExecutor는 추가적인 스레드 간 핸드오버 오버헤드를 피해요. 콜백은 보통 결과를 output 버퍼에 추가하는 ResultFuture에 넘길 뿐이에요. 거기서부터 레코드 발행과 체크포인트 부기와의 상호작용을 포함한 무거운 로직은 어차피 전용 스레드 풀에서 일어나요.

DirectExecutor는 org.apache.flink.util.concurrent.Executors.directExecutor() 또는 com.google.common.util.concurrent.MoreExecutors.directExecutor()로 얻을 수 있어요.

참고: 이는 Java API에만 적용돼요. Python API에서는 비동기 결과를 대기하면 돼요.

주의사항 (Caveats)

AsyncFunction은 다중 스레드로 호출되지 않음 (The AsyncFunction is not called Multi-Threaded)

명시적으로 지적하고 싶은 흔한 혼동은 AsyncFunction이 다중 스레드 방식으로 호출되지 않는다는 것이에요. AsyncFunction의 인스턴스는 하나만 존재하며 스트림의 해당 파티션의 각 레코드에 대해 순차적으로 호출돼요. asyncInvoke(...) 메서드가 빠르게 반환하고 (클라이언트의) 콜백에 의존하지 않으면 적절한 비동기 I/O가 되지 않아요. 예를 들어 다음 패턴은 차단하는 asyncInvoke(...) 함수를 만들어 비동기 동작을 무효화해요:

  • 결과를 다시 받을 때까지 조회/쿼리 메서드 호출이 차단되는 데이터베이스 클라이언트 사용
  • asyncInvoke(...) 메서드 안에서 비동기 클라이언트가 반환한 future-타입 객체를 차단/대기

AsyncFunction(AsyncWaitOperator)은 작업 그래프 어디에서나 사용할 수 있지만 SourceFunction/SourceStreamTask에 체이닝될 수는 없어요.

재시도가 활성화되면 더 큰 큐 용량이 필요할 수 있음 (May Need Larger Queue Capacity If Retry Enabled)

새 재시도 기능은 더 큰 큐 용량 요구를 초래할 수 있으며, 최대 수는 대략 다음과 같이 평가할 수 있어요:

inputRate * retryRate * avgRetryDuration

예를 들어 inputRate = 100 records/sec인 작업에서 요소의 1%가 평균적으로 1회 재시도를 트리거하고 평균 재시도 시간이 60초라면 추가 큐 용량 요구는 다음과 같아요:

100 records/sec * 1% * 60s = 60

즉 작업 큐에 60만큼 더 추가하는 것이 unordered 출력 모드에서 처리량에 영향을 주지 않을 수 있어요. ordered 모드에서는 head 요소가 핵심이며, 그 요소가 미완료로 오래 유지될수록 연산자가 제공하는 처리 지연이 길어져요. 같은 타임아웃 제약으로 실제로 더 많은 재시도를 얻는다면 재시도 기능은 head 요소의 미완료 시간을 늘릴 수 있어요.

큐 용량이 커지면(백프레셔를 완화하는 흔한 방법) OOM 위험이 커져요. 사실 ListState 저장의 경우 이론적 상한은 Integer.MAX_VALUE이므로 큐 용량의 한계도 같지만, 프로덕션에서 큐 용량을 너무 크게 늘릴 수는 없어요. 작업 병렬도를 늘리는 것이 더 실행 가능한 방법일 수 있어요.

더 알아보기 (Learn more)