뭐라고요 제가 3년차라고요 데이터 엔지니어 2025 회고
·
주간 · 월간 회고
2025 키워드끊임없는 의사결정의 연속AI를 어떻게 하면 더 잘 사용할 수 있을 것인가2025년 프로젝트1. 데이터 파이프라인 리팩토링리팩토링의 목표는 두 가지였다.하드코딩 된 설정들을 전부 밖으로 빼서 Airflow Variables와 Connection에 넘겨 관리하자각 테이블별, 단계별로 독립적인 Task로 만들어서 실패와 성공 이력을 관리하고 백필을 용이하게 만들자 결론적으로 두 목표는 모두 이뤘고, 추가 개발 사항과 리팩토링으로 인해 생기는 또 다른 문제를 해결했다. 1-1. 추가 개발 사항 [실시간 데이터의 수요]실시간으로 데이터를 모니터링할 수 없냐는 요청 사항이 있었다. 운영 DB의 안전성 문제로 Kafka 커넥션을 맺기에는 어려웠고, 기존 방식 대로 쿼리하여 데이터를 가져왔어야 했다. 데..
[Spark] java.lang.illegalArgumentException: Too large frame 에러 원인 및 해결 방법
·
Data Engineering/Spark
1. 개요대용량 과거 데이터를 적재하는데 해당 에러 발생 2. 장애 현상해당 에러로 인해서 계속 계속 Spark 연결 시도 그러나 데이터 프레임을 불러오지는 못함 3. 원인 분석불러오려는 데이터가 너무 커서 Spark 가 감당 못한다는 에러이다. 확인해 보니 잘 불러왔던 데이터에 비해 row수가 3배 이상 컸다. 대부분은 collect()를 사용해서 파티션 단위의 데이터들을 드라이버로 모을 때 부하가 일어나 생긴다는데, 지금 ETL에서는 collect는 사용하고 있지 않았다. 이상 없는 데이터들과 비교했을 때 지나치게 row 수가 많다는 것. 이게 원인인 것 같다. 4. 해결 방법원래는 1년치씩 데이터를 불러왔는데 해당 데이터만 6개월로 끊어 적재했다.
대용량 데이터 Airflow로 백필(backfill) 시, java.io.IOException: 장치에 남은 공간이 없음
·
Data Engineering/Airflow
1. 개요HIVE로 테이블을 적재하는데, 모든 데이터에 파티션 데이터를 추가하다 보니 스키마를 완전히 바꿔야 했다. 그로 인해 불가피하게 재적재해야 하는 상황 발생. 기존에는 쉘 명령어로 재적재를 하였지만, 이번에는 Airflow backfill 기능을 사용해서 재적재하기로 했다. 왜냐하면 어느 부분이 성공하고 실패했는지 확인하기 쉽고, 검증 단계까지 포함되어 있기 때문이다. 2. 장애 현상Task에 해당 에러 발생주목했던 에러는 장치에 남은 공간이 없다는 에러였다.INFO - py4j.protocol.Py4JJavaError: An error occurred while calling None.org.apache.spark.api.java.JavaSparkContext.INFO - : java.io.IO..
Spark window 함수 예제 (공부용)
·
Study/SQL
1) 사용자별 최신 이벤트 1건만 남기기목표: user_id마다 가장 최근 event_time 1행만 남긴다.WITH t AS ( SELECT user_id, -- 사용자 식별자 event_time, -- 이벤트 발생 시간 payload, -- 부가 데이터(원하는 컬럼) ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY event_time DESC) AS rn -- PARTITION BY user_id: 사용자별로 그룹을 나눔 -- ORDER BY event_time DESC: 각 사용자 그룹 안에서 최신 시간이 먼저 오게 정렬 -- ROW_NUMBER(): 정렬된 순서대로 1,2..
Iceberg vs Hudi vs DeltaLake -> Iceberg 선택한 이유
·
Data Engineering/Infra
Hive를 최대한 걷어내고 이를 대체할 테이블 포맷 세 가지를 비교했다.현재 환경실시간 데이터보다는 일반 배치 데이터가 주를 이룬다. 그렇다고 실시간 데이터를 아예 배제하지는 않는다. 추후 실시간 데이터를 적재했을 때의 상황도 고려해야 한다. 아키텍처를 새로 구성한다고 해서 한순간 하나를 완전히 걷어낼 수는 없다. spark와 trino 두 개를 모두 사용해야 되는 상황이다. 비교공식 Docs도 찾아보고 실사용자들의 후기를 많이 읽어보고 공부한 결과 현재 데이터 배치와 상황을 고려했을 때 Iceberg가 가장 적절했다.항목IceberghudiDeltaLake호환성Spark, Trino 등등 다양한 SQL 엔진과의 호환성 뛰어남Spark Structured Streaming, Kafka 실시간 데이터 처리..
Rag Vector DB로 Milvus를 선택한 이유
·
Data Engineering/Infra
문제 상황사내 문서를 학습한 AI 모델의 수요가 생겼다. 당장 파인 튜닝을 진행할 수는 없고, 쉽고 빠르게 구현할 수 있는 RAG를 적용하기로 했다.RAG를 적용하려면 문서를 벡터화해서 저장해 놓을 벡터 DB가 필요하다. 해결 과정Vector DB 고려 사항 1. 확장성내가 아키텍처를 구성할 때 가장 먼저 생각하는 것 중 하나이다. 예상했던 것보다 더 많은 데이터가 들어온다면? 노드를 더 추가해야 된다면? 새로운 기능이 구현될 때 추가하기 용이한가? 등등 추가 작업이 있을 때 편한지를 먼저 생각한다.2. 권한 부여모든 사람이 같은 권한을 가지고 접속하면 안 되기 때문에 사람마다 다른 권한을 부여할 수 있는지 체크해야 한다.3. 온프레미스매니지드 DB보다 상황에 맞춰 A부터 Z까지 환경 설정할 수 있는 온..