AWS Managed Service for Apache Flink(구 Kinesis Data Analytics) 환경에서 PyFlink를 활용해 Kinesis Data Streams(KDS)의 실시간 데이터를 집계하는 파이프라인을 구축하던 중 발생한 일련의 오류와 해결 과정을 기록한다.
로컬 환경(Python 3.11)에서 apache-flink==1.15.0 설치 시도 중 subprocess-exited-with-error 발생.
apache-beam 등 의존성 라이브러리가 Python 3.11의 빌드 구조와 충돌함.코드 실행 시 ModuleNotFoundError: No module named 'flink' 및 생성자 호출 에러 발생.
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)
CREATE TABLE 구문에서 'connector' = 'kinesis'를 사용했으나, Flink 엔진이 커넥터를 찾지 못함.
NoSuchMethodError 등 런타임 에러 발생 가능성이 매우 높음.flink-sql-connector-kinesis-4.2.0-1.18.jar 파일로 교체.
S3에 코드를 올리고 실행했으나 RestHandlerException 발생 및 데이터 수신 불가.
kinesis.analytics.flink.run.optionspython / Value: app.pyjarfile / Value: flink-sql-connector-kinesis-4.2.0-1.18.jarkinesis:DescribeStream, kinesis:GetRecords, kinesis:PutRecord 등 추가.
Reference