| 일 | 월 | 화 | 수 | 목 | 금 | 토 |
|---|---|---|---|---|---|---|
| 1 | 2 | 3 | 4 | 5 | ||
| 6 | 7 | 8 | 9 | 10 | 11 | 12 |
| 13 | 14 | 15 | 16 | 17 | 18 | 19 |
| 20 | 21 | 22 | 23 | 24 | 25 | 26 |
| 27 | 28 | 29 | 30 |
- 컨테이너 삭제
- Kafka
- dag 작성
- docker hub
- ETL
- dag
- Hive
- 데이터레이크
- airflow.cfg
- selenium
- airflow
- yarn
- Django
- spark
- Serializer
- truncate
- SQL
- ELT
- AWS
- redshift
- docker
- 웹 스크래핑
- Django Rest Framework(DRF)
- 웹 크롤링
- 데이터마트
- snowflake
- docker-compose
- 데이터파이프라인
- 데이터 웨어하우스
- 알고리즘
- Today
- Total
목록전체 글 (70)
개발 기록장
배경우리 팀에서는 Airflow를 목적에 따라 3개의 EC2 환경, 3개의 Airflow를 사용하고 있었다. 3개로 나누어 사용했던 가장 큰 이유는 EC2 인스턴스의 크기, 스케줄링 타이밍의 문제였다. DAG이 동시에 돌아가게 되면 OOM이 발생하기 때문이다.이원화된 airflow는 관리 포인트가 많았다. 첫번째는 Disk Full 상황이 자주 발생한다는 점이었다. Disk Full인지 주기적으로 ssh로 각 인스턴스에 접근하여 꾸준히 df를 통해 디스크를 확인하고 비워주는 작업을 했다. 그런데 점점 DAG이 늘어날 수록 이런 상황은 더 자주 발생했고, Airflow가 멈추는 사태도 더 많아졌다. 두번째, Airflow가 멈추면 다시 정상화 하고, 데이터를 Backfill하는데 많은 리소스가 들었다. 우..
Airflow Dag에서 task를 다룰 때, 병렬로 task를 처리하면 더 효율적인 경우가 있다. 나의 CASE에서는 API 호출을 통해 각 task에서 데이터를 받아오고, 데이터를 마지막에 merge해야 했다. 만약 직렬로 처리한다면 task들은 순차적으로 API를 통해 데이터를 받아오게 되고, 비슷한 task가 여러 개로 늘어나면 merge까지 오랜 시간이 걸린다. API 호출은 CPU 연산이 아니라 대기시간(I/O wait)이 대부분이기 때문에 병렬 실행이 더 효율적이다. 또 마지막에 merge task에서 앞의 task가 하나라도 제대로 작동하지 않았다면, merge가 실행되지 않도록 하여 데이터 정합성을 해치지 않도록 제어하기 수월하다.(물론 직렬 실행에서도 관리할 수 있긴 하다. -> Tri..
Cloud Run상황 설명목표는 GCP 리소스들을 이용하여 CDC 파이프라인을 만드는 것이었다. Background는 GCP A 계정에서 생성된 Cloud SQL for postgresql에서 생성되는 데이터를 GCP B계정의 DW(Bigquery)로 저장하는 CDC 파이프라인을 만드는 것이다. 예상 파이프라인 아키텍처는 Cloud SQL for postgresql(A 계정) -> Datasteam(B 계정) -> GCS(B 계정) -> Pub/Sub(B 계정) -> Dataflow(B 계정) -> Bigquery(B 계정)이었다. 또한 요구 사항은 A 계정의 Cloud SQL for postgresql은 Cloud sql proxy로 접근해야했다.(다른 계정의 cloud sql for postgresq..
AWS S3 스토리지로 데이터를 적재할 때의 문제점S3는 서비스 DB로 적합하지 않음지하철 데이터는 15초의 주기로 생성되고, 이를 반영해 대시보드에서도 실시간 성을 유지해야 한다. 그러나 S3에 적재하게 되면, Kafka Topic으로부터 데이터를 저장하는 것뿐만 아니라 데이터를 대시보드에서 출력하는 데에도 시간이 오래 걸린다. 또 실시간 지하철 정보 데이터는 축적되어 저장할 필요가 없기 때문에 용량이 커다란 스토리지가 필요하지 않다.ELT의 필요성API로부터 받아온 데이터에는 대시보드 시각화에 필요 없는 정보와 보기 편하도록 변환해야 하는 값들도 존재한다. 기존에는 이 값들을 태블로 대시보드 시각화 과정에서 정리하려 했으나, 태블로에서 이 값들을 처리하는 과정에서도 시간이 오래 걸린다.(실시간 성을 ..
기존에는 Kafka에서 Producer.py와 Consumer.py를 이용하여 데이터를 처리했다. 이 방식은 간단한 데이터를 처리할 경우에는 직관적으로 빠르게 코드를 작성하여 처리할 수 있다는 장점이 있지만, 확장성/ 모니터링의 측면에서는 적합하지 않다는 단점이 있었다.그래서 우리는 Kafka의 Connector의 사용을 고려하기 시작했다. 우리가 프로젝트에서 받아와야할 데이터는 서울 열린데이터 광장에서 지하철 데이터를 실시간으로 받아와야하고, API의 호출 횟수 제한이 있으므로 모니터링이 굉장히 중요한 부분이었다. 또 커넥터를 사용한다면 변경 사항이 있을 때, 따로 코드의 작성 없이 Kafka 상에서 커넥터 설정만을 수정해 사용할 수 있으므로 간편하다고 생각했다.Kafka ConnectorKafka C..
상황: 실시간 지하철 데이터를 Kafka를 이용하여 처리해아함먼저 강의에서 다루었던 데이터 전달 방식인 Producer.py와 Consumer.py를 작성하여 Kafka 환경을 Test했다.데이터는 서울 열린 데이터 광장의 노선별 실시간 지하철 데이터를 이용: https://data.seoul.go.kr/dataList/OA-12601/A/1/datasetView.doProducer.py: Topic을 생성하고 데이터를 Topic으로 전송한다.노선별 지하철 정보를 받기위해 subway 리스트를 만들어 반복문을 돌려 API 요청을 넣었음Consumer.py에서 노선별 지하철 정보 저장을 위해, 실시간 지하철 데이터와 함께 지하철 노선명도 함께 Topic을 통해 전달함생성된 Topic 이름: subway_r..
데이터 전송에는 벌크 형과 스트리밍 형 두 종류가 있다.객체 스토리지와 데이터 수집: 분산 스토리지에 데이터 읽어들이기빅데이터는 대부분 확장성이 높은 분산 스토리지(distributed storage)에 저장됨기본적으로 대랑으로 파일을 저장하기위한 객체 스토리지(object storage)가 많이 사용됨Hadoop이라면 HDFS, 클라우드 서비스라면 Amazon S3 등객체 스토리지의 내부 처리에는 다수의 물리적 서버와 하드디스크가 존재하므로 일부의 HW가 고장나더라도 데이터 손실X데이터의 읽기, 쓰기를 다수의 하드웨어에 분산 -> 데이터의 양이 늘어나도 성능 유지객체 스토리지 구조는 데이터 양이 많을 때는 우수하지만, 소량의 데이터를 처리할 때는 오버헤드가 너무 크므로 비효율적데이터 수집: 수집한 데이..
Queue, enumerate문제 설명운영체제의 역할 중 하나는 컴퓨터 시스템의 자원을 효율적으로 관리하는 것입니다. 이 문제에서는 운영체제가 다음 규칙에 따라 프로세스를 관리할 경우 특정 프로세스가 몇 번째로 실행되는지 알아내면 됩니다.1. 실행 대기 큐(Queue)에서 대기중인 프로세스 하나를 꺼냅니다.2. 큐에 대기중인 프로세스 중 우선순위가 더 높은 프로세스가 있다면 방금 꺼낸 프로세스를 다시 큐에 넣습니다.3. 만약 그런 프로세스가 없다면 방금 꺼낸 프로세스를 실행합니다. 3.1 한 번 실행한 프로세스는 다시 큐에 넣지 않고 그대로 종료됩니다.예를 들어 프로세스 4개 [A, B, C, D]가 순서대로 실행 대기 큐에 들어있고, 우선순위가 [2, 1, 3, 2]라면 [C, D, A, B] 순으로..