1、kafka重要参数acks值对比:0、1和all(-1)

acks值 数据可靠性 延迟性 适用场景 备注
0 最低(可能丢失数据) 最低 对延迟敏感,允许数据丢失(如日志、指标) 生产者不等待服务器响应
1 中等(leader写入即确认) 中等 默认配置,平衡可靠性与延迟(如一般消息队列) 需leader副本持久化,follower未同步时可能丢失
all/-1 最高(所有ISR副本同步) 最高 要求强一致性(如金融交易)

需配置min.insync.replicas控制最小同步副本数

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


 

Logo

腾讯云面向开发者汇聚海量精品云计算使用和开发经验,营造开放的云计算技术生态圈。

更多推荐