프로젝트 진행 시, 전처리된 csv를 s3로 보내는 task에 넘기려고 하였다. 하지만 여기서 오류가 발생했다.
round_df = round_df.to_dict()
round_df = json.dumps(round_df)
return round_df
# S3에 저장하는 함수
def save_to_s3(data, s3_bucket, s3_key, aws_access_key_id, aws_secret_access_key):
result_all = json.loads(data)
result_all = pd.DataFrame(result_all)
s3 = boto3.client('s3',
aws_access_key_id=aws_access_key_id,
aws_secret_access_key=aws_secret_access_key)
s3.put_object(Body=result_all.to_csv(encoding='utf-8'), Bucket=s3_bucket, Key=s3_key)#변경
이런식으로 csv를 json으로 바꾸어서 전송하였다. Airflow의 XCom을 통해서 task간 output을 전달하는데, XCom은 일반적으로 Python 객체를 지원하지만 JSON과 같은 직렬화된 형식으로 데이터를 저장하고 전달한다고 한다.
물론 이 부분은 airflow를 더 공부해서 깊게 파보아야할 것 같다.