Airflow에서 하나의 Dag는 다수의 operator 혹은 task로 구성됨
start_date 는 2/5일 이면 실행 날짜는 2/6일 이다.
위 대그의 경우 execution date는 2/5일이 된다.
schedule interval이 "30****"
매시 30분 마다 실행된다
매일 0시 30분
"30 0 ***"
일요일 매시 30분
"30 *** 0"
0은 일요일 6은 토요일
schedule 0****이라면 -> 매시 0분에 한번씩 실행 hourly dag
start date = 2021 2/4 00:00:00
그럼 한시에 실행된다!
redshift에서 큰 데이터를 테이블로 복사하는 방식 중 올바른 것은
복사할 레코드들을 파일로 저장해서 s3로 올린 후에 거기서 reshift로 벌크 복사
distinct로 중복 제거를 하게 되면, 완벽하게 값이 같지 않다면 실행되지 않는 경우가 발생
전형적인 데이터 인프라를 만들 때 해야할 일
프로덕션 데이터베이스 -> 데이터 웨어하우스
구현하려는 ETL 소개
production MySQL tables OLTP
를 aws redshift OLAP로
##ETL
mysql에서 테이블을 읽어서 redshift로 Insert into
근데 mysql에서 레코드가 많다면 insert into하다가 시간이 많이 흘러버림
소스를 데이터를 가져와서 file로 저장해서
s3에 업로드 해서
redshift에 copy로 bulk update
airflow dag에서 s3 접근
iam user를 만들고 s3 버킷에대한 읽기 쓰기 권한 설정하고 access key 와 secret key를 사용
redshift s3접근
redshift에 s3를 접근할 수 있는 역할을 만들고 이를 redshift에 지정
Access Key ID와 Secret Access Key를 사용하는 걸로 바뀜
루트 사용자의 키들을 사용하면 해킹시 AWS 자원들을 마음대로 사용가능 -> 여러번 사고가 남
우리가 사용할 best practice는
iam을 통해서 별도의 사용자를 만들고
그사용자에게 해당 s3 bucket을 읽고 쓸 수 있는 권한 제공
사용자의 ACcess Key ID와 Secret Access key 를 사용
이 키도 주도적으로 변경해서 해킹이 될 경우의 피해를 최소화
iam 가서 직접 권한 연결 후에
{
"Version": "2012-10-17", "Statement": [
{
"Effect": "Allow", "Action": [
"s3:GetBucketLocation",
"s3:ListAllMyBuckets" ],
"Resource": "arn:aws:s3:::*" },
{
"Effect": "Allow", "Action": "s3:*", "Resource": [
"arn:aws:s3:::grepp-data-engineering",
"arn:aws:s3:::grepp-data-engineering/*" ]
} ]
}
이런식으로 정책 생성후 연결해서 key들을 가지고 연결하면 된다
airflow 내부적으로 판다스로 가져오게 되어있음 그레서 index와 header를 가져올 것인지에 대한 옵션을 설정해줘야함.
프로덕션 디비 테이블의 incremental update 방식
deleted : 레코드를 삭제하지 않고 deleted 를 true
● Daily Update이고 테이블의 이름이 A이고 MySQL에서 읽어온다면
● ROW_NUMBER로 직접 구현하는 경우
○ 먼저 Redshift의 A 테이블의 내용을 temp_A로 복사
○ MySQL의 A 테이블의 레코드 중 modified의 날짜가 지난 일(execution_date)에 해당하는 모든
레코드를 읽어다가 temp_A로 복사
■ 아래는 MySQL에 보내는 쿼리. 결과를 파일로 저장한 후 S3로 업로드하고 COPY 수행
● SELECT * FROM A WHERE DATE(modified) = DATE(execution_date)
○ temp_A의 레코드들을 primary key를 기준으로 파티션한 다음에 modified 값을 기준으로
DESC 정렬해서, 일련번호가 1인 것들만 다시 A로 복사
● Daily Update이고 테이블의 이름이 A이고 MySQL에서 읽어온다면
● S3ToRedshiftOperator로 구현하는 경우
○ query 파라미터로 아래를 지정
● SELECT * FROM A WHERE DATE(modified) = DATE(execution_date)
○ method 파라미터로 “UPSERT”를 지정
○ upsert_keys 파라미터로 Primary key를 지정
■ 앞서 nps 테이블이라면 “id” 필드를 사용
일단은 mysql을 해결했다. mysql 문제를 해결하니 redshift에 테이블이 없는 문제가 발생
그래서 아무거나 만들어줬더니 바보 같이 컬럼을 맞추지 못해 오류가 생겨서 mysql의 테이블을 미리 조회후에 필드를 만들고 성공
mysql은 어떻게 해결을 했는지 모르겠다. 너무 여러방법을 써서...
1) sudo 로 airflow에 mysqlclient를 설치
2) airflow connection에 mysql 설치
“conn_type”: “mysql”, “client”: “mysql-connector-python”}
근데 또 v2 같은 경우는 ui에서는 안되고 터미널에서는 된다 이유는 모르겠다....
airflow에서 추천하는 방식으로 Incremental Update를 구현했다면 Backfill이 쉬워짐
하지만 이를 어떻게 실행하나?
하루씩 31번 실행?
execution date를 이용해서 구현할 수 는 있다.
한번에 여러 날짜를 동시에 실행 가능한가?
구현 방법에 따라 한번에 하나씩 실행하는 것이 안전할 수 있음
이를 제어해주는 DAG 파라미터가 max_active_runs
일반적으로 production db에서는 읽기 전용으로 slave를 하나 구성한다.
한번에 읽으면 생길 수 있는 문제
ex) mysql to redshift 같은 경우
max_active_runs 파라미터를 1로 지정하고
backfill을 할 때, 한번에 한날짜만 실행이 되도록 지정해 둠
airflow dags backfill dag_id -s 2018-07-01 -e 2018-08-01
● This assumes the followings:
○ catchUp이 True로 설정되어 있음
○ execution_date을 사용해서 Incremental update가 구현되어 있음
dag 구성이 위와 같이 구현되어있어야함!!
만약 full refresh 라면 시간과 자원의 낭비
● start_date부터 시작하지만 end_date은 포함하지 않음
● 실행순서는 날짜/시간순은 아니고 랜덤. 만일 날짜순으로 하고 싶다면
DAG default_args의 depends_on_past를 True로 설정 default_args = {
'depends_on_past': True,
먼저 모든 DAG가 backfill을 필요로 하지는 않음 ○ Full Refresh를 한다면 backfill은 의미가 없음
● 여기서 backfill은 일별 혹은 시간별로 업데이트하는 경우를 의미함
○ 마지막 업데이트 시간 기준 backfill을 하는 경우라면 (Data Warehouse 테이블에 기록된 시간
기준) 이런 경우에도 execution_date을 이용한 backfill은 필요하지 않음
● 데이터의 크기가 굉장히 커지면 backfill 기능을 구현해 두는 것이 필수
○ airflow가 큰 도움이 됨
○ 하지만 데이터 소스의 도움 없이는 불가능
어떻게 backfill로 구현할 것인가
○ 제일 중요한 것은 데이터 소스가 backfill 방식을 지원해야함
○ “execution_date”을 사용해서 업데이트할 데이터 결정
○ “catchup” 필드를 True로 설정
○ start_date/end_date을 backfill하려는 날짜로 설정
○ 다음으로 중요한 것은 DAG 구현이 execution_date을 고려해야 하는 것이고 idempotent 해야함
데이터 파이프라인을 만들 때, backfill과 upsert를 생각하면서 코드를 작성해야할 듯
ETL 데이터 시스템 밖에서 안으로 데이터를 가져와서
데이터 시스템 내에서 join하여 새로운 데이터를 만들어 내는 것이 ELT
Airflow 장점
데이터 파이프라인을 세밀하게 제어 가능
다양한 데이터 소스와 데이터 웨어하우스를 지원
백필이 쉬움
start_date
execution_date
catchup
스케일링 방식
클라우드 버전, k8s
○ 데이터 파이프라인에 관한 정보를 수집하는 것이 중요
○ 결국 데이터 카탈로그가 필요
● 데이터 품질 체크
○ 입력 데이터와 출력 데이터
● 코드 실패를 어설프게 복구하려는 것보다는 깔끔하게 실패하는 것이 좋음
● 가능하면 Full Refresh
○ Incremental Update를 쓸 수 밖에 없다면 Backfill 방식을 먼저 생각해둘 것 -> Airflow가 필요한
이유
● 주기적인 청소 (데이터, 테이블, Dag)