안녕하세요 강사님 dbt랑 airflow를 현재 수강중인 직장인입니다. 배울수록 굉장히 활용범위가 넓은 툴이라고 생각이 됩니다. 두가지 질문 드리고 싶은데요 1) 현재 저는 dbt+Airflow 기반으로 CRM 분석 마트 테이블을 팀에 적용하려고 하고 있습니다. 현재 raw테이블을 자동화해서 airflow로 db에 적재하여 사용하고 있는데요, 조인 결합, 가공을 통한 2차, 3차 테이블들은 현재 수동으로 만들어지고 있고 이를 자동화하려고 하고 있는데 현재 운영 및 성과 분석을 위한 너무 많은 테이블이 생기면서 점점 복잡해지고 있어 처음 만든 저도 헷갈려지는 단계에 왔는데...설계 관리(테이블간 관계, 단계로직을 정리하여 적재)하는 것에 대한 노하우가 있으신지 궁금합니다. 그리고 설계 이후 dbt로 만든 모델을 팀원들(분석가나 마케터) 쉽게 활용할 수 있게 하려면, 어떤 방식으로 문서화나 공유를 하시나요? 2) Airflow DAG에서 dbt run/test를 통합할 때, 실행 단위를 모델 단위로 쪼개는 게 좋은가요, 아니면 전체 프로젝트 단위로 돌리는 게 좋은가요?
안녕하세요, 강의 잘 들었습니다. 아직 실무에 flink를 사용해 본 적이 없어 조금 더 구체적인 장점등을 알고 싶어 질문 드리게 되었습니다. 사실 기존에는 kafka만을 사용해서 실시간 데이터 처리를 하는 사례도 많았던 것 같은데 이 때 추가로 flink를 사용할 때 kafka만을 사용할 때 보다 어떤 부분이 더 나은지 등에 대해 조금 궁금해져서 질문 드립니다.
좋은 강의를 제공해주셔서 감사합니다. Kafka를 처음 설정할 때 Partition 개수를 1개로 두고 추후 확장하는 방식이 일반적인가요? 아니면 초기에 적절한 Partition 개수를 여유롭게 미리 설정하는 것이 좋을까요? 또한, 이러한 설정은 멀티 노드 Kafka 클러스터 구성 여부에 따라 달라질 수 있는지도 궁금합니다.
안녕하세요. 강의 수강 중 궁금한게 있어서 질문 남깁니다. 강의는 VM 에 Ubuntu 와 kafka 를 설치하는 것으로 진행되는데, Docker 를 사용하는 것과 VM 을 사용하는 것에 차이가 있나요? VM 이 아니라 Docker 로 Kafka 를 띄우거나, Ubuntu 를 띄우고 Kafka 를 설치해도 동일하지 않나 생각이 들더라고요.
안녕하세요 선생님 클러스터 설정 시 오류가 발생하여 질문 드립니다. ㅜㅜ 다른 질문 글들을 참고하여 Cluster 1 삭제 후 재설치도 해보았고, putty로 접속하여 rm -rf /dfs/nn 명령어로 디렉토리 삭제 후 cluster 재설치도 해보았는데 계속 오류가 발생합니다. 원인과 해결 방법이 있을지 문의 드립니다.. * stderr로그 일부 Caused by: org.apache.hadoop.ipc.RemoteException(org.apache.hadoop.hdfs.server.namenode.SafeModeException): Cannot create directory /tmp. Name node is in safe mode. The reported blocks 0 has reached the threshold 0.9990 of total blocks 0. The number of live datanodes 0 needs an additional 1 live datanodes to reach the minimum number 1. Safe mode will be turned off automatically once the thresholds have been reached. NamenodeHostName:server01.hadoop.com at org.apache.hadoop.hdfs.server.namenode.FSNamesystem.newSafemodeException(FSNamesystem.java:1448) at org.apache.hadoop.hdfs.server.namenode.FSNamesystem.checkNameNodeSafeMode(FSNamesystem.java:1435) at org.apache.hadoop.hdfs.server.namenode.FSNamesystem.mkdirs(FSNamesystem.java:3100) at org.apache.hadoop.hdfs.server.namenode.NameNodeRpcServer.mkdirs(NameNodeRpcServer.java:1123) at org.apache.hadoop.hdfs.protocolPB.ClientNamenodeProtocolServerSideTranslatorPB.mkdirs(ClientNamenodeProtocolServerSideTranslatorPB.java:696) at org.apache.hadoop.hdfs.protocol.proto.ClientNamenodeProtocolProtos$ClientNamenodeProtocol$2.callBlockingMethod(ClientNamenodeProtocolProtos.java) at org.apache.hadoop.ipc.ProtobufRpcEngine$Server$ProtoBufRpcInvoker.call(ProtobufRpcEngine.java:523) at org.apache.hadoop.ipc.RPC$Server.call(RPC.java:991) at org.apache.hadoop.ipc.Server$RpcCall.run(Server.java:869) at org.apache.hadoop.ipc.Server$RpcCall.run(Server.java:815) at java.security.AccessController.doPrivileged(Native Method) at javax.security.auth.Subject.doAs(Subject.java:422) at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1875) at org.apache.hadoop.ipc.Server$Handler.run(Server.java:2675) at org.apache.hadoop.ipc.Client.getRpcResponse(Client.java:1499) at org.apache.hadoop.ipc.Client.call(Client.java:1445) at org.apache.hadoop.ipc.Client.call(Client.java:1355) at org.apache.hadoop.ipc.ProtobufRpcEngine$Invoker.invoke(ProtobufRpcEngine.java:228) at org.apache.hadoop.ipc.ProtobufRpcEngine$Invoker.invoke(ProtobufRpcEngine.java:116) at com.sun.proxy.$Proxy9.mkdirs(Unknown Source) at org.apache.hadoop.hdfs.protocolPB.ClientNamenodeProtocolTranslatorPB.mkdirs(ClientNamenodeProtocolTranslatorPB.java:640) at sun.reflect.GeneratedMethodAccessor8.invoke(Unknown Source) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.lang.reflect.Method.invoke(Method.java:498) at org.apache.hadoop.io.retry.RetryInvocationHandler.invokeMethod(RetryInvocationHandler.java:422) at org.apache.hadoop.io.retry.RetryInvocationHandler$Call.invokeMethod(RetryInvocationHandler.java:165) at org.apache.hadoop.io.retry.RetryInvocationHandler$Call.invoke(RetryInvocationHandler.java:157) at org.apache.hadoop.io.retry.RetryInvocationHandler$Call.invokeOnce(RetryInvocationHandler.java:95) at org.apache.hadoop.io.retry.RetryInvocationHandler.invoke(RetryInvocationHandler.java:359) at com.sun.proxy.$Proxy10.mkdirs(Unknown Source) at org.apache.hadoop.hdfs.DFSClient.primitiveMkdir(DFSClient.java:2339) ... 18 more
6강에서 Cursor의 AI를 이용해서 프로그램의 뼈대를 만드는 부분에서 강사님과 같은 프롬프트를 써도 미묘하게 폴더 구조가 강사님과 다르게 작성이 됩니다. 프롬프트가 같다면 무조건 결과가 같은건지 아니면 다를 수도 있는지.. 가능하다면 강의에서 구성한 뼈대 부분을 파일로 제공해 주실 수 있는지 문의드립니다. (전부 화면을 보고 다 따라치려고 해도 강의 화면만 보고는 그것도 무리인거 같습니다.) 기왕이면 이후의 남은 강의들에도 소스 코드를 제공해 주시면 좋을 것 같습니다.
안녕하세요 CLI로 Trigger 기능을 수행하는 부분 강의를 듣던 중에, Web UI에서 Trigger를 누르면 정상적으로 수행 되지만, 쉘 스크립트 커맨드로 airflow dags trigger <DAG 이름>이라는 명령어를 실행했을 때 아래와 같은 실패 로그가 나타나서 문의드립니다. 혹시 커맨드라인으로 실행하면 {{data_interval_end}} 와 같은 템플릿을 적용할 수 없나요?
안녕하세요. 섹션 12 강의를 듣는중인데 airflow 디렉토리 밑에 custom_image 디렉토리가 이미 하나 있어야 하더라구요. 그런데 제 airflow 디렉토리 밑에는 해당 디렉토리가 없습니다. 여태 진행한 강의는 분명 빠짐없이 들었는데 제가 실수로 놓친 부분이 있는 것 같습니다.. 다시 찾아 듣고자 하는데 어느 강의인지 찾지를 못하고 있습니다. 죄송하지만 혹시 해당 부분 몇 강에서 진행하셨는지 알 수 있을까요?
안녕하세요 선생님, 현재 데이터 엔지니어 직무 면접을 준비하고 있는 수강생입니다. 저는 이번 면접에서 ETL 아키텍처 및 데이터 파이프라인 구성 과 관련된 주제를 중심으로 준비하고 있습니다. 수업을 통해 많은 내용을 배우고 있지만, 솔직히 말씀드리면 양이 너무 방대하다 보니 모든 부분을 다 소화하기가 쉽지 않습니다. 시간도 빠듯해 점점 불안감이 커지고 있습니다. 그래서 정말 간절한 마음으로 도움을 부탁드립니다. 혹시 면접에서 특히 중요한 핵심 개념 이나 자주 나오는 질문 , 그리고 실제 사례 중심으로 정리하면 좋은 부분 이 무엇인지 방향을 잡아주실 수 있을까요? 혹은 강의만 따로 추천 해주시면 감사하겠습니다. 제가 부족한 부분을 더 보완해서 면접에서 꼭 좋은 결과를 내고 싶습니다. 늘 열정적으로 가르쳐주시는 것에 감사드리며, 간절히 도움을 청합니다. 감사합니다.
안녕하세요 강사님. 좋은 강의 감사드립니다. 항상 강의 잘듣고 있습니다 😀 다름이 아니라, 복합키 디코딩 관련 질문이 있습니다. 아래와 같이 Group By와 Window Session 함수를 결합한 CTAS절입니다. CREATE OR REPLACE TABLE MASTER WITH ( KAFKA_TOPIC = 'master', KEY_FORMAT = 'JSON', VALUE_FORMAT = 'JSON' ) AS SELECT TRID AS KEY, AS_VALUE(TRID) AS "trid", WINDOWSTART AS "min_time", WINDOWEND AS "max_time", (WINDOWEND - WINDOWSTART) AS "duration", MIN timestamp ) AS "@timestamp", COLLECT_LIST(service) AS services, COLLECT_SET(system) AS systems FROM ORIGINAL_STREAM WINDOW SESSION (5 SECONDS) GROUP BY TRID EMIT CHANGES; 제가 기대한 값으로는 master라는 토픽의 key에 trid와 windowstart 값으로 결합된 JSON 형식의 값이 저장되는 것이었습니다. ksqldb에서 print 문으로 topic을 조회하면 잘 읽히지만, kafka-consumer에서 topic을 조회하면, 디코딩 부분에서 깨져서 조회가 됩니다. 명령어: ./kafka/bin/ kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic cpm_master --from-beginning --property print.key=true 현재는 총 두 개의 쿼리를 추가적으로 사용하여 id 값을 컨슈머가 읽을 수 있도록 정제하고 있습니다. 혹시 이 문제에 대해서 아신다면 답변 주시면 감사하겠습니다!
echo $CONFLUENT_HOME이 정확한 위치인것 확인했고, zookeeper.properties와 zookeeper-server-start가 정확한 위치인 것도 확인했습니다. 그러나 zookeeper-server-start $CONFLUENT_HOME/etc/kafka/zookeeper.properties를 치면 아무것도 나오지 않습니다. 에러 문구조차도 없네요
안녕하세요. Cooperative Sticky Rebalancing할 때 컨슈머 3만 reassign rebalance된다고 하셨습니다. 궁금한건 컨슈머 1 2에서 topic A, B / Partion 1 2 각각 총 4개는 유지된다고 하셨는데 이떄 poll은 계속 동작하고 있는건지 궁금합니다. 즉, 3이 죽고 이걸 새로 붙이는 컨슈머 1 2에 있는 모든 토픽과 모든 파티션에 poll이 정지되나요? 아니면 같은 토픽 컨슈머만 정지되나요? (여기선 A P1이 있어고 A에 P3가 붙으므로 P1이 P3를 위한 assign rebalance 때문에 컨슈머 1 A 토픽 poll 전체 정지) reassign rebalance: 부분 할당 rebalance: 전체 할당 으로 이해하면 되는지도 궁금합니다. 감사합니다.
안녕하세요. Consumer Fetcher관련 주요 파라미터와 Fetcher 메커니즘의 이해 강의 5:03에서 망냑 파티셩니 10개 있으면 최대 10MB 가져올 수 있따 이렇게 말씀하셨는데 컨슈머를 띄울 때 파티션별로 각 서버마다 따로 뜨게 하시는지 한 컨슈머 서버에 여러 파티션을 구독하게 띄우시는지 궁금합니다. 이를 설정하는 기준이 있으신지도 궁금합니다. 제가 카프카를 처음 공부하고 있어서 혹시 질문이 잘못되었다면 알려주시면 감사하겠습니다. 감사합니다.