标签:orm points loop oca build cal api res dex
近日,Hudi社区合并了 Flink 引擎的基础实现(HUDI-1327),这意味着 Hudi 开始支持 Flink 引擎。
当前 Flink 版本的 Hudi 只支持读取 Kafka 数据,sink到 COW 类型的 Hudi 表中,其他功能还在完善。
这里我们简要介绍下如何从 Kafka 读取数据 写出到 Hudi表。
hadoop-2.7.2 kafka_2.11-2.1.0 flink-1.11.2 hudi-0.7.0
官方地址:https://hudi.apache.org/ gihub地址:https://github.com/apache/hudi
git clone https://github.com/apache/hudi.git && cd hudi
修改hudi/pom.xml
1> 切换release-0.7.0分支,修改pom.xml中hadoop对应集群版本
2> window环境下注释掉<module>hudi-integ-test</module>和<module>packaging/hudi-integ-test-bundle</module>
3> flink-1.11.2对应parquet版本为1.11.1,可修改pom.xml到对应版本(若flink/lib下存在其他本版以该版本保持一致),本次测试不修改保留1.10.1版本
mvn clean package -DskipTests
hoodie.datasource.write.recordkey.field=uuid hoodie.datasource.write.partitionpath.field=ts bootstrap.servers=hadoop01:9092,hadoop02:9092,hadoop03:9092 hoodie.deltastreamer.keygen.timebased.timestamp.type=EPOCHMILLISECONDS hoodie.deltastreamer.keygen.timebased.output.dateformat=yyyy/MM/dd hoodie.datasource.write.keygenerator.class=org.apache.hudi.keygen.TimestampBasedAvroKeyGenerator hoodie.embed.timeline.server=false hoodie.deltastreamer.schemaprovider.source.schema.file=hdfs:///tmpdir/hudi/test/config/flink/schema.avsc hoodie.deltastreamer.schemaprovider.target.schema.file=hdfs:///tmpdir/hudi/test/config/flink/schema.avsc
{ "type":"record", "name":"stock_ticks", "fields":[{ "name": "uuid", "type": "string" }, { "name": "ts", "type": "long" }, { "name": "symbol", "type": "string" },{ "name": "year", "type": "int" },{ "name": "month", "type": "int" },{ "name": "high", "type": "double" },{ "name": "low", "type": "double" },{ "name": "key", "type": "string" },{ "name": "close", "type": "double" }, { "name": "open", "type": "double" }, { "name": "day", "type":"string" } ]}
sudo -u hdfs hadoop fs -mkdir -p /tmpdir/hudi/test/config/flink hadoop fs -put schema.avsc /tmpdir/hudi/test/config/flink/ hadoop fs -put hudi-conf.properties /tmpdir/hudi/test/config/flink/
#创建主题 /opt/apps/kafka/bin/kafka-topics.sh --create --zookeeper hadoop01:2181,hadoop02:2181,hadoop03:2181/kafka --replication-factor 2 --partitions 3 --topic mytest #查看主题 /opt/apps/kafka/bin/kafka-topics.sh --list --zookeeper hadoop01:2181,hadoop02:2181,hadoop03:2181/kafka #生产者(测试) /opt/apps/kafka/bin/kafka-console-producer.sh --broker-list hadoop01:9092,hadoop02:9092,hadoop03:9092 --topic mytest
测试数据
{"close":0.27172637588467297,"day":"2","high":0.4493211149337879,"key":"840ef1","low":0.030714155934507215,"month":11,"open":0.7762668153935262,"symbol":"77c-40d6-8412-6859d4757727","ts":1608361070161,"uuid":"840ef1ff-b77c-40d6-8412-6859d4757727","year":120}
/opt/flink-1.11.2/bin/flink run -c org.apache.hudi.HoodieFlinkStreamer -m yarn-cluster -d -yjm 2048 -ytm 3096 -ys 4 -ynm hudi_on_flink_test -p 1 -yD env.java.opts=" -XX:+TraceClassPaths -XX:+TraceClassLoading" hudi-flink-bundle_2.11-0.7.0.jar --kafka-topic mytest --kafka-group-id hudi_on_flink --kafka-bootstrap-servers hadoop01:9092,hadoop02:9092,hadoop03:9092 --table-type COPY_ON_WRITE --target-base-path hdfs:///tmpdir/hudi/test/data/hudi_on_flink \ --target-table hudi_on_flink --props hdfs:///tmpdir/hudi/test/config/flink/hudi-conf.properties --checkpoint-interval 3000 \ --flink-checkpoint-path hdfs:///flink/hudi/hudi_on_flink_cp
HoodieFlinkStreamer参数介绍
参数 | 描述 |
---|---|
--kafka-topic | 必选,kafka source主题 |
--kafka-group-id | 必选,kafka消费者组 |
--kafka-bootstrap-servers | 必选,kafka bootstrap.servers 如 node1:9092,node2:9092,node3:9092 |
--flink-checkpoint-path | 可选,flink checkpoint 路径 |
--flink-block-retry-times | 可选,默认10,当最近的hoodie instant未完成时重试的次数 |
--flink-block-retry-interval | 可选,默认1,当最近的hoodie instant未完成时,两次尝试之间的秒数 |
--target-base-path | 必选,目标hoodie表的基本路径(如果路径不存在将被创建) |
--target-table | 必选,Hive中目标表的名称 |
--table-type | 必选,表的类型。COPY_ON_WRITE 或 MERGE_ON_READ |
--props | 可选,属性配置文件的路径(local或hdfs)。有hoodie客户端、schema提供者、键生成器和数据源的配置。 |
--hoodie-conf | 可选,可以在属性文件中设置的任何配置(使用参数"--props")也可以使用此参数传递命令行。 |
--source-ordering-field | 可选,以决定如何打破输入数据中具有相同键的记录之间的联系字段。默认值:‘ts‘保存记录的unix时间戳 |
--payload-class | 可选,HoodieRecordPayload的子类,它在GenericRecord上工作。实现你自己的,如果你想做一些事情而不是覆盖现有的值 |
--op | 可选,接受以下值之一:UPSERT(默认),INSERT(当输入纯粹是新数据/插入时使用,以提高速度) |
--filter-dupes | 可选,默认false。是否应该在插入/大容量插入之前删除/过滤源中的重复记录 |
--commit-on-errors | 可选,默认false。即使某些记录写入失败也要提交 |
--checkpoint-interval | 可选,默认5000毫秒。Flink checkpoint间隔 |
新搭建的测试环境未做过jar包的兼容性整合。job能正常启动后,当消费消息write_process时会发生如下异常:
java.io.IOException: Could not perform checkpoint 6 for operator write_process (1/1). at org.apache.flink.streaming.runtime.tasks.StreamTask.triggerCheckpointOnBarrier(StreamTask.java:892) at org.apache.flink.streaming.runtime.io.CheckpointBarrierHandler.notifyCheckpoint(CheckpointBarrierHandler.java:113) at org.apache.flink.streaming.runtime.io.CheckpointBarrierAligner.processBarrier(CheckpointBarrierAligner.java:137) at org.apache.flink.streaming.runtime.io.CheckpointedInputGate.pollNext(CheckpointedInputGate.java:93) at org.apache.flink.streaming.runtime.io.StreamTaskNetworkInput.emitNext(StreamTaskNetworkInput.java:158) at org.apache.flink.streaming.runtime.io.StreamOneInputProcessor.processInput(StreamOneInputProcessor.java:67) at org.apache.flink.streaming.runtime.tasks.StreamTask.processInput(StreamTask.java:351) at org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.runMailboxStep(MailboxProcessor.java:191) at org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.runMailboxLoop(MailboxProcessor.java:181) at org.apache.flink.streaming.runtime.tasks.StreamTask.runMailboxLoop(StreamTask.java:566) at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:536) at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:721) at org.apache.flink.runtime.taskmanager.Task.run(Task.java:546) at java.lang.Thread.run(Thread.java:748) Caused by: org.apache.flink.runtime.checkpoint.CheckpointException: Could not complete snapshot 6 for operator write_process (1/1). Failure reason: Checkpoint was declined. at org.apache.flink.streaming.api.operators.StreamOperatorStateHandler.snapshotState(StreamOperatorStateHandler.java:215) at org.apache.flink.streaming.api.operators.StreamOperatorStateHandler.snapshotState(StreamOperatorStateHandler.java:156) at org.apache.flink.streaming.api.operators.AbstractStreamOperator.snapshotState(AbstractStreamOperator.java:314) at org.apache.flink.streaming.runtime.tasks.SubtaskCheckpointCoordinatorImpl.checkpointStreamOperator(SubtaskCheckpointCoordinatorImpl.java:614) at org.apache.flink.streaming.runtime.tasks.SubtaskCheckpointCoordinatorImpl.buildOperatorSnapshotFutures(SubtaskCheckpointCoordinatorImpl.java:540) at org.apache.flink.streaming.runtime.tasks.SubtaskCheckpointCoordinatorImpl.takeSnapshotSync(SubtaskCheckpointCoordinatorImpl.java:507) at org.apache.flink.streaming.runtime.tasks.SubtaskCheckpointCoordinatorImpl.checkpointState(SubtaskCheckpointCoordinatorImpl.java:266) at org.apache.flink.streaming.runtime.tasks.StreamTask.lambda$performCheckpoint$8(StreamTask.java:921) at org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$1.runThrowing(StreamTaskActionExecutor.java:47) at org.apache.flink.streaming.runtime.tasks.StreamTask.performCheckpoint(StreamTask.java:911) at org.apache.flink.streaming.runtime.tasks.StreamTask.triggerCheckpointOnBarrier(StreamTask.java:879) ... 13 more Caused by: org.apache.hudi.exception.HoodieUpsertException: Failed to upsert for commit time 20210125212048 at org.apache.hudi.table.action.commit.AbstractWriteHelper.write(AbstractWriteHelper.java:62) at org.apache.hudi.table.action.commit.FlinkUpsertCommitActionExecutor.execute(FlinkUpsertCommitActionExecutor.java:47) at org.apache.hudi.table.HoodieFlinkCopyOnWriteTable.upsert(HoodieFlinkCopyOnWriteTable.java:66) at org.apache.hudi.table.HoodieFlinkCopyOnWriteTable.upsert(HoodieFlinkCopyOnWriteTable.java:58) at org.apache.hudi.client.HoodieFlinkWriteClient.upsert(HoodieFlinkWriteClient.java:110) at org.apache.hudi.operator.KeyedWriteProcessFunction.snapshotState(KeyedWriteProcessFunction.java:121) at org.apache.flink.streaming.util.functions.StreamingFunctionUtils.trySnapshotFunctionState(StreamingFunctionUtils.java:120) at org.apache.flink.streaming.util.functions.StreamingFunctionUtils.snapshotFunctionState(StreamingFunctionUtils.java:101) at org.apache.flink.streaming.api.operators.AbstractUdfStreamOperator.snapshotState(AbstractUdfStreamOperator.java:90) at org.apache.hudi.operator.KeyedWriteProcessOperator.snapshotState(KeyedWriteProcessOperator.java:58) at org.apache.flink.streaming.api.operators.StreamOperatorStateHandler.snapshotState(StreamOperatorStateHandler.java:186) ... 23 more Caused by: java.lang.RuntimeException: org.apache.hudi.exception.HoodieException: org.apache.hudi.exception.HoodieException: java.util.concurrent.ExecutionException: java.lang.NoClassDefFoundError: org/apache/hudi/avro/HoodieAvroWriteSupport at org.apache.hudi.client.utils.LazyIterableIterator.next(LazyIterableIterator.java:121) at java.util.Iterator.forEachRemaining(Iterator.java:116) at org.apache.hudi.table.action.commit.BaseFlinkCommitActionExecutor.lambda$execute$0(BaseFlinkCommitActionExecutor.java:120) at java.util.LinkedHashMap.forEach(LinkedHashMap.java:684) at org.apache.hudi.table.action.commit.BaseFlinkCommitActionExecutor.execute(BaseFlinkCommitActionExecutor.java:118) at org.apache.hudi.table.action.commit.BaseFlinkCommitActionExecutor.execute(BaseFlinkCommitActionExecutor.java:68) at org.apache.hudi.table.action.commit.AbstractWriteHelper.write(AbstractWriteHelper.java:55) ... 33 more Caused by: org.apache.hudi.exception.HoodieException: org.apache.hudi.exception.HoodieException: java.util.concurrent.ExecutionException: java.lang.NoClassDefFoundError: org/apache/hudi/avro/HoodieAvroWriteSupport at org.apache.hudi.execution.FlinkLazyInsertIterable.computeNext(FlinkLazyInsertIterable.java:73) at org.apache.hudi.execution.FlinkLazyInsertIterable.computeNext(FlinkLazyInsertIterable.java:38) at org.apache.hudi.client.utils.LazyIterableIterator.next(LazyIterableIterator.java:119) ... 39 more Caused by: org.apache.hudi.exception.HoodieException: java.util.concurrent.ExecutionException: java.lang.NoClassDefFoundError: org/apache/hudi/avro/HoodieAvroWriteSupport at org.apache.hudi.common.util.queue.BoundedInMemoryExecutor.execute(BoundedInMemoryExecutor.java:143) at org.apache.hudi.execution.FlinkLazyInsertIterable.computeNext(FlinkLazyInsertIterable.java:69) ... 41 more Caused by: java.util.concurrent.ExecutionException: java.lang.NoClassDefFoundError: org/apache/hudi/avro/HoodieAvroWriteSupport at java.util.concurrent.FutureTask.report(FutureTask.java:122) at java.util.concurrent.FutureTask.get(FutureTask.java:192) at org.apache.hudi.common.util.queue.BoundedInMemoryExecutor.execute(BoundedInMemoryExecutor.java:141) ... 42 more Caused by: java.lang.NoClassDefFoundError: org/apache/hudi/avro/HoodieAvroWriteSupport at org.apache.hudi.io.storage.HoodieFileWriterFactory.newParquetFileWriter(HoodieFileWriterFactory.java:59) at org.apache.hudi.io.storage.HoodieFileWriterFactory.getFileWriter(HoodieFileWriterFactory.java:47) at org.apache.hudi.io.HoodieCreateHandle.<init>(HoodieCreateHandle.java:85) at org.apache.hudi.io.HoodieCreateHandle.<init>(HoodieCreateHandle.java:66) at org.apache.hudi.io.CreateHandleFactory.create(CreateHandleFactory.java:34) at org.apache.hudi.execution.CopyOnWriteInsertHandler.consumeOneRecord(CopyOnWriteInsertHandler.java:83) at org.apache.hudi.execution.CopyOnWriteInsertHandler.consumeOneRecord(CopyOnWriteInsertHandler.java:40) at org.apache.hudi.common.util.queue.BoundedInMemoryQueueConsumer.consume(BoundedInMemoryQueueConsumer.java:37) at org.apache.hudi.common.util.queue.BoundedInMemoryExecutor.lambda$null$2(BoundedInMemoryExecutor.java:121) at java.util.concurrent.FutureTask.run(FutureTask.java:266) at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) ... 1 more
发现存在该类,于是猜测与集群中有jar冲突
通过 find /opt -name "*hudi*.jar" 查找。并未发现有相关的jar(曾切换到不同版本的hadoop集群测试结果一致异常)
结论:不是jar冲突引起。
添加 -yD env.java.opts="-XX:+TraceClassLoading -XX:+TraceClassPaths"参数,查看HoodieAvroWriteSupport 没有被加载。于是开始定位异常代码位置,猜测为静态代码中初始化该类时就已经异常导致,将 new AvroSchemaConverter().convert(schema)写到外面来定位具体的缺失类
Caused by: java.lang.ClassNotFoundException: org.apache.parquet.schema.Type
mac
到编译的本地maven依赖仓库找出以下parquet相关jar拷贝到flink/lib下尝试解决
Caused by: java.lang.NoClassDefFoundError: org/apache/parquet/hadoop/ParquetInputFormat
at org.apache.parquet.HadoopReadOptions$Builder.<init>(HadoopReadOptions.java:87)
查看classloadtrace信息已加载HadoopReadOptions来确认parquet-hadoop包已经被加载,所以确定ParquetInputFormat类(来自parquet-hadoop)的异常由其他缺失造成,于是查看HadoopReadOptions类中引入的类及ParquetInputFormat继承的父类等。最终发现ParquetInputFormat继承的父类org.apache.hadoop.mapreduce.lib.input.FileInputFormat
(所属hadoop-mapreduce-client-coreye-x.x.x.jar)并未被加载。
重新提交任务。若发现有如下异常,删除hdfs表目录下的文件重新提交(由于之前步骤的异常造成的损坏)。
java.io.IOException: Could not perform checkpoint 3 for operator instant_generator (1/1). at org.apache.flink.streaming.runtime.tasks.StreamTask.triggerCheckpointOnBarrier(StreamTask.java:892) at org.apache.flink.streaming.runtime.io.CheckpointBarrierHandler.notifyCheckpoint(CheckpointBarrierHandler.java:113) at org.apache.flink.streaming.runtime.io.CheckpointBarrierAligner.processBarrier(CheckpointBarrierAligner.java:137) at org.apache.flink.streaming.runtime.io.CheckpointedInputGate.pollNext(CheckpointedInputGate.java:93) at org.apache.flink.streaming.runtime.io.StreamTaskNetworkInput.emitNext(StreamTaskNetworkInput.java:158) at org.apache.flink.streaming.runtime.io.StreamOneInputProcessor.processInput(StreamOneInputProcessor.java:67) at org.apache.flink.streaming.runtime.tasks.StreamTask.processInput(StreamTask.java:351) at org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.runMailboxStep(MailboxProcessor.java:191) at org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.runMailboxLoop(MailboxProcessor.java:181) at org.apache.flink.streaming.runtime.tasks.StreamTask.runMailboxLoop(StreamTask.java:566) at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:536) at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:721) at org.apache.flink.runtime.taskmanager.Task.run(Task.java:546) at java.lang.Thread.run(Thread.java:748) Caused by: java.lang.InterruptedException: Last instant costs more than 10 second, stop task now at org.apache.hudi.operator.InstantGenerateOperator.doCheck(InstantGenerateOperator.java:199) at org.apache.hudi.operator.InstantGenerateOperator.prepareSnapshotPreBarrier(InstantGenerateOperator.java:119) at org.apache.flink.streaming.runtime.tasks.OperatorChain.prepareSnapshotPreBarrier(OperatorChain.java:266) at org.apache.flink.streaming.runtime.tasks.SubtaskCheckpointCoordinatorImpl.checkpointState(SubtaskCheckpointCoordinatorImpl.java:249) at org.apache.flink.streaming.runtime.tasks.StreamTask.lambda$performCheckpoint$8(StreamTask.java:921) at org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$1.runThrowing(StreamTaskActionExecutor.java:47) at org.apache.flink.streaming.runtime.tasks.StreamTask.performCheckpoint(StreamTask.java:911) at org.apache.flink.streaming.runtime.tasks.StreamTask.triggerCheckpointOnBarrier(StreamTask.java:879) ... 13 more
上述问题基本是由于集群没有相关jar包导致,根据上面的问题定位方法能解决大部分在遇到NoClassDefFoundError、ClassNotFoundException等异常问题。
Hudi on flink v0.7.0 使用遇到的问题及解决办法
标签:orm points loop oca build cal api res dex
原文地址:https://www.cnblogs.com/felixzh/p/14478700.html