数仓开发flink相关概念
1、kafka重要参数acks值对比:0、1和all(-1)
| acks值 | 数据可靠性 | 延迟性 | 适用场景 | 备注 |
|---|---|---|---|---|
| 0 | 最低(可能丢失数据) | 最低 | 对延迟敏感,允许数据丢失(如日志、指标) | 生产者不等待服务器响应 |
| 1 | 中等(leader写入即确认) | 中等 | 默认配置,平衡可靠性与延迟(如一般消息队列) | 需leader副本持久化,follower未同步时可能丢失 |
| all/-1 | 最高(所有ISR副本同步) | 最高 | 要求强一致性(如金融交易) |
需配置 |
2、在实时开发中 Flink Java DataStream API 与 Flink SQL 在关键维度的差异:
| 维度 | Flink Java DataStream API | Flink SQL |
|---|---|---|
| 底层执行性能 | 无抽象层,调优空间大(新手易写差) | 优化器自动生成最优执行计划(新手也能写出高效逻辑) |
| 开发效率 | 低(需写大量模板代码) | 高(一行 SQL 搞定核心逻辑,无模板代码) |
| 灵活性 | 极致灵活(支持任意复杂逻辑) | 中等(标准化场景最优,复杂逻辑需 UDF) |
| 运维成本 | 高(改逻辑需重打包 / 重部署) | 低(改 SQL 脚本即可,无需打包) |
| 新手友好性 | 低(需掌握 Java+Flink 核心 API) | 高(SQL 语义通用,无需写代码) |
| 标准化 ETL 适配度 | 可用但没必要 | 最优选择 |
3、使用Flink Java DataStream API开发时的调用逻辑
| 阶段 | 代码位置 | 核心操作 | 调用关系 |
|---|---|---|---|
| 作业入口 | main(String[] args) |
程序启动,接收命令行参数 | JVM 启动 Flink 作业,入口方法 |
| 参数解析与校验 | main 方法内前半部分 |
解析参数→校验必传项→类型转换→合法性检查 | 主方法直接执行,无外部调用 |
| Kerberos 全局配置 | main 方法内 |
设置系统级 krb5.conf 路径 |
主方法直接执行,消费者/生产者共用此配置 |
| Flink 环境初始化 | main 方法内 |
初始化执行环境→配置 Checkpoint/StateBackend/并行度/重启策略 | 主方法调用 Flink API,构建作业运行上下文 |
| Kafka 配置构建 | buildKafkaProps(pt, isConsumer) |
根据isConsumer参数,生成消费者/生产者配置 |
主方法调用2次:① isConsumer=true → 给 Kafka Source 用 ② isConsumer=false → 给 Kafka Sink 用 |
| 算子链定义(懒加载) | main 方法内算子链式调用 |
addSource()→map()→addSink(),定义数据流转路径 |
主方法调用 Flink 算子 API,仅定义执行计划,不处理数据 |
| 作业提交 | env.execute("作业名") |
提交执行计划到 Flink 集群,启动作业 | 主方法触发作业执行,这是懒执行的开关 |
| 数据处理循环(作业运行时) | TrajectoryMapper.map(String value) |
每条 Kafka 数据触发一次map调用,处理空值/解析 JSON/填充默认值 |
Flink Runtime 调用,每条数据独立处理 |
///spark任务提交
/opt/spark-3.4.2-bin-hadoop3/bin/spark-shell \
--master yarn \
--queue u_xa_css_batch \
--driver-memory 20G \
--num-executors 50 \
--executor-cores 4 \
--executor-memory 20G \
--conf spark.driver.maxResultSize=32G \
--conf spark.yarn.executor.memoryOverhead=8G \
--conf spark.default.parallelism=1000 \
--conf spark.sql.shuffle.partitions=300 \
--conf spark.shuffle.memoryFraction=0.3 \
--conf spark.driver.extraJavaOptions=-XX:+UseG1GC \
--conf spark.executor.extraJavaOptions=-XX:+UseG1GC \
--conf spark.executor.userClassPathFirst=true \
--conf spark.driver.userClassPathFirst=true \
--conf spark.speculation=true
///flink任务提交
/opt/flink-1.18.1/bin/flink run-application \
-t yarn-application \
-Dsecurity.kerberos.login.use-ticket-cache=false \
-Dsecurity.kerberos.login.keytab=/home/
-Dsecurity.kerberos.login.principal=u@HADOOP.COM \
-Dsecurity.kerberos.login.contexts=Client,KafkaClient \
-Djobmanager.memory.process.size=64G \
-Dtaskmanager.memory.process.size=32G \
-Dtaskmanager.memory.network.fraction=0.25 \
-Dtaskmanager.memory.managed.size=0G \
-Dtaskmanager.numberOfTaskSlots=8 \
-Dyarn.application.queue=u_xa_css_stream \
-Dyarn.application.name=xa_css_TargetMain \
-plugins" \
-c lbs.task.MainNoTargetAggLink \
/home/flink/f1k-rzfx-1.0-SNAPSHOT.jar \
///////////////////////////
/opt/flink-1.18.1/bin/flink run-application \
-t yarn-application \
-Dsecurity.kerberos.login.use-ticket-cache=false \
-Dsecurity.kerberos.login.keytab=/opt/kerberos/user.keytab \
-Dsecurity.kerberos.login.principal=user@HADOOP.COM \
-Dsecurity.kerberos.login.contexts=Client,KafkaClient \
-Djobmanager.memory.process.size=8G \
-Dyarn.applicationmaster.vcores=1
-Dyarn.containers.vcores=4 \
-Dyarn.containers.max=5 \
-Dtaskmanager.memory.process.size=16G \
-Dtaskmanager.memory.network.fraction=0.1 \
taskmanager.memory.managed.fraction=0.2 \
-Dtaskmanager.numberOfTaskSlots=4 \
-Dyarn.application.queue=test_stream \
-Dyarn.application.name=yym-test-Job \
-Dyarn.provided.lib.dirs="/user/yym/lib:/user/yym/plugins" \
/opt/flink-1.18.1/jars/flink-kafka-trajectory-1.0-SNAPSHOT-jar-with-dependencies.jar \
-c com.trajectory.FlinkSqlTrajectoryJob \
--kafka.bootstrap.servers=kafka1:9092,kafka2:9092 \
--flink.checkpoint.interval=5000 \
--flink.parallelism=20 \
--flink.checkpoint.dir=hdfs:///flink/checkpoints/trajectory-job \
--kafka.consumer.group.id=test_group \
--kafka.input.topic=trajectory-input \
--kafka.output.topic=trajectory-output \
--kafka.krb5.conf.path=/etc/krb5.conf \
--default.track.id=default1 \
--default.longitude=11 \
--default.latitude=22
更多推荐
所有评论(0)