[Troubleshooting] PyFlink 1.18 + AWS Managed Flink + KDS 연동

정성헌·2026년 4월 16일

1. 개요

AWS Managed Service for Apache Flink(구 Kinesis Data Analytics) 환경에서 PyFlink를 활용해 Kinesis Data Streams(KDS)의 실시간 데이터를 집계하는 파이프라인을 구축하던 중 발생한 일련의 오류와 해결 과정을 기록한다.


2. 문제 상황 1: Python 버전과 라이브러리 호환성

로컬 환경(Python 3.11)에서 apache-flink==1.15.0 설치 시도 중 subprocess-exited-with-error 발생.

원인 분석

  • 팩트: Apache Flink 1.15는 Python 3.10 이상의 환경을 공식 지원하지 않음. 특히 apache-beam 등 의존성 라이브러리가 Python 3.11의 빌드 구조와 충돌함.
  • 해결: Python 3.11 환경을 유지하기 위해 Flink 버전을 1.18.0으로 업그레이드하여 설치 진행.

3. 문제 상황 2: 모듈 참조 및 객체 생성 오류

코드 실행 시 ModuleNotFoundError: No module named 'flink' 및 생성자 호출 에러 발생.

원인 분석

  • Import 경로: 패키지명은 apache-flink이지만, 실제 코드 내 임포트 경로는 pyflink임. from flink.table은 존재하지 않는 경로.
  • 객체 생성: TableEnvironment는 생성자를 직접 호출하지 않고 .create() 정적 메서드를 사용해야 함.

교정된 코드

from pyflink.table import EnvironmentSettings, TableEnvironment

# 잘못된 방식: t_env = TableEnvironment(setting)
# 올바른 방식:
setting = EnvironmentSettings.new_instance().in_streaming_mode().build()
t_env = TableEnvironment.create(setting)

4. 문제 상황 3: Kinesis Connector 인식 불가 (JAR 버전 mismatch)

CREATE TABLE 구문에서 'connector' = 'kinesis'를 사용했으나, Flink 엔진이 커넥터를 찾지 못함.

원인 분석

  • 팩트 1: Flink SQL을 쓰기 위해서는 일반 Connector가 아닌 SQL 전용 통합 JAR(Fat Jar)가 필요함.
  • 팩트 2: Flink 1.18 엔진에 1.15용 JAR를 사용하면 클래스 로딩 시 NoSuchMethodError 등 런타임 에러 발생 가능성이 매우 높음.
  • 해결: flink-sql-connector-kinesis-4.2.0-1.18.jar 파일로 교체.

5. 문제 상황 4: AWS 환경에서의 배포 및 권한 문제

S3에 코드를 올리고 실행했으나 RestHandlerException 발생 및 데이터 수신 불가.

원인 분석

  • EntryPoint 미지정: AWS Managed Flink는 ZIP 내의 어떤 파일이 메인인지 모름. '런타임 속성(Runtime Properties)' 설정 필수.
  • IAM 권한 결여: 기본 생성된 역할에는 Kinesis 읽기/쓰기 권한이 빠져 있었음.

최종 설정 데이터

  1. Runtime Properties:
    • Group ID: kinesis.analytics.flink.run.options
    • Key: python / Value: app.py
    • Key: jarfile / Value: flink-sql-connector-kinesis-4.2.0-1.18.jar
  2. IAM Policy: kinesis:DescribeStream, kinesis:GetRecords, kinesis:PutRecord 등 추가.


6. 결론 및 회고

  1. 버전 정합성: PyFlink는 Python, Flink 엔진, Connector JAR의 3박자 버전이 완벽히 맞아야 한다.
  2. 환경의 차이: 로컬(Windows)과 배포(AWS Linux) 환경의 경로 체계 차이를 인지하고 하드코딩을 피해야 한다.
  3. 권한 확인: "시동이 안 걸리면 설정 문제, 달리다 멈추면 권한 문제"라는 가설을 데이터(CloudWatch)로 증명하는 과정이 중요하다.

Reference

  • Apache Flink Documentation (v1.18)
  • Maven Repository: flink-sql-connector-kinesis
profile
develop myself

0개의 댓글