从EMR向MSK写入数据时遇TimeoutException:元数据无指定Topic
问题描述
我在Amazon EMR(版本emr-6.10.0)上运行一个Scala编写的Spark应用,尝试通过IAM认证方式将数据写入Amazon MSK(Kafka版本3.4.0)。
我已通过以下命令在Amazon MSK中创建了名为hm.motor.avro的Topic:
bin/kafka-topics.sh \ --bootstrap-server=b-1.myemr.xxx.c12.kafka.us-west-2.amazonaws.com:9098,b-2.myemr.xxx.c12.kafka.us-west-2.amazonaws.com:9098 \ --command-config=config/client.properties \ --create \ --topic=hm.motor.avro \ --partitions=3 \ --replication-factor=2
以下是使用IAM认证向MSK写入数据的相关Spark代码:
val query = df.writeStream .format("kafka") .option( "kafka.bootstrap.servers", "b-1.myemr.xxx.c12.kafka.us-west-2.amazonaws.com:9098,b-2.myemr.xxx.c12.kafka.us-west-2.amazonaws.com:9098", ) .option("kafka.security.protocol", "SASL_SSL") .option("kafka.sasl.mechanism", "AWS_MSK_IAM") .option("kafka.sasl.jaas.config", "software.amazon.msk.auth.iam.IAMLoginModule required;") .option("kafka.sasl.client.callback.handler.class", "software.amazon.msk.auth.iam.IAMClientCallbackHandler") .option("topic", "hm.motor.avro") .option("checkpointLocation", "/tmp/checkpoint") .start()
build.sbt(我使用的Spark版本为3.3.1,与Amazon EMR内置版本一致)
name := "IngestFromS3ToKafka" version := "1.0" scalaVersion := "2.12.17" resolvers += "confluent" at "https://packages.confluent.io/maven/" val sparkVersion = "3.3.1" libraryDependencies ++= Seq( "org.apache.spark" %% "spark-core" % sparkVersion % "provided", "org.apache.spark" %% "spark-sql" % sparkVersion % "provided", "org.apache.hadoop" % "hadoop-common" % "3.3.3" % "provided", "org.apache.hadoop" % "hadoop-aws" % "3.3.3" % "provided", "com.amazonaws" % "aws-java-sdk-bundle" % "1.12.397" % "provided", "org.apache.spark" %% "spark-avro" % sparkVersion, "org.apache.spark" %% "spark-sql-kafka-0-10" % sparkVersion, "io.delta" %% "delta-core" % "2.2.0", "za.co.absa" %% "abris" % "6.3.0", "software.amazon.msk" % "aws-msk-iam-auth" % "1.1.6" ) ThisBuild / assemblyMergeStrategy := { case PathList("module-info.class") => MergeStrategy.discard case x if x.endsWith("/module-info.class") => MergeStrategy.discard case PathList("org", "apache", "spark", "unused", "UnusedStubClass.class") => MergeStrategy.first case x if x.contains("io.netty.versions.properties") => MergeStrategy.discard case x => val oldStrategy = (ThisBuild / assemblyMergeStrategy).value oldStrategy(x) }
应用通过sbt assembly打包,使用Java 1.8(EMR默认Java版本)。
但无论以YARN cluster还是client模式执行spark-submit,均出现如下错误:
Caused by: org.apache.spark.SparkException: Writing job aborted at org.apache.spark.sql.errors.QueryExecutionErrors$.writingJobAbortedError(QueryExecutionErrors.scala:767) at org.apache.spark.sql.execution.datasources.v2.V2TableWriteExec.writeWithV2(WriteToDataSourceV2Exec.scala:409) at org.apache.spark.sql.execution.datasources.v2.V2TableWriteExec.writeWithV2$(WriteToDataSourceV2Exec.scala:353) at org.apache.spark.sql.execution.datasources.v2.WriteToDataSourceV2Exec.writeWithV2(WriteToDataSourceV2Exec.scala:302) at org.apache.spark.sql.execution.datasources.v2.WriteToDataSourceV2Exec.run(WriteToDataSourceV2Exec.scala:313) at org.apache.spark.sql.execution.datasources.v2.V2CommandExec.result$lzycompute(V2CommandExec.scala:43) at org.apache.spark.sql.execution.datasources.v2.V2CommandExec.result(V2CommandExec.scala:43) at org.apache.spark.sql.execution.datasources.v2.V2CommandExec.executeCollect(V2CommandExec.scala:49) at org.apache.spark.sql.Dataset.collectFromPlan(Dataset.scala:3932) at org.apache.spark.sql.Dataset.$anonfun$collect$1(Dataset.scala:3161) at org.apache.spark.sql.Dataset.$anonfun$withAction$2(Dataset.scala:3922) at org.apache.spark.sql.execution.QueryExecution$.withInternalError(QueryExecution.scala:554) at org.apache.spark.sql.Dataset.$anonfun$withAction$1(Dataset.scala:3920) at org.apache.spark.sql.catalyst.QueryPlanningTracker$.withTracker(QueryPlanningTracker.scala:107) at org.apache.spark.sql.execution.SQLExecution$.withTracker(SQLExecution.scala:224) at org.apache.spark.sql.execution.SQLExecution$.executeQuery$1(SQLExecution.scala:114) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$7(SQLExecution.scala:139) at org.apache.spark.sql.catalyst.QueryPlanningTracker$.withTracker(QueryPlanningTracker.scala:107) at org.apache.spark.sql.execution.SQLExecution$.withTracker(SQLExecution.scala:224) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$6(SQLExecution.scala:139) at org.apache.spark.sql.execution.SQLExecution$.withSQLConfPropagated(SQLExecution.scala:245) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$1(SQLExecution.scala:138) at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:779) at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:68) at org.apache.spark.sql.Dataset.withAction(Dataset.scala:3920) at org.apache.spark.sql.Dataset.collect(Dataset.scala:3161) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runBatch$17(MicroBatchExecution.scala:669) at org.apache.spark.sql.catalyst.QueryPlanningTracker$.withTracker(QueryPlanningTracker.scala:107) at org.apache.spark.sql.execution.SQLExecution$.withTracker(SQLExecution.scala:224) at org.apache.spark.sql.execution.SQLExecution$.executeQuery$1(SQLExecution.scala:114) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$7(SQLExecution.scala:139) at org.apache.spark.sql.catalyst.QueryPlanningTracker$.withTracker(QueryPlanningTracker.scala:107) at org.apache.spark.sql.execution.SQLExecution$.withTracker(SQLExecution.scala:224) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$6(SQLExecution.scala:139) at org.apache.spark.sql.execution.SQLExecution$.withSQLConfPropagated(SQLExecution.scala:245) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$1(SQLExecution.scala:138) at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:779) at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:68) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runBatch$16(MicroBatchExecution.scala:664) at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken(ProgressReporter.scala:375) at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken$(ProgressReporter.scala:373) at org.apache.spark.sql.execution.streaming.StreamExecution.reportTimeTaken(StreamExecution.scala:68) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.runBatch(MicroBatchExecution.scala:664) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runActivatedStream$2(MicroBatchExecution.scala:256) at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken(ProgressReporter.scala:375) at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken$(ProgressReporter.scala:373) at org.apache.spark.sql.execution.streaming.StreamExecution.reportTimeTaken(StreamExecution.scala:68) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runActivatedStream$1(MicroBatchExecution.scala:219) at org.apache.spark.sql.execution.streaming.ProcessingTimeExecutor.execute(TriggerExecutor.scala:67) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.runActivatedStream(MicroBatchExecution.scala:213) at org.apache.spark.sql.execution.streaming.StreamExecution.$anonfun$runStream$1(StreamExecution.scala:307) at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:779) at org.apache.spark.sql.execution.streaming.StreamExecution.org$apache$spark$sql$execution$streaming$StreamExecution$$runStream(StreamExecution.scala:285) ... 1 more Caused by: org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 6.0 failed 4 times, most recent failure: Lost task 0.3 in stage 6.0 (TID 106) (ip-xxx-xxx-xxx-xxx.xxx.com executor 1): org.apache.kafka.common.errors.TimeoutException: Topic hm.motor.avro not present in metadata after 60000 ms.
注意最后一行错误:
org.apache.kafka.common.errors.TimeoutException: Topic hm.motor.avro not present in metadata after 60000 ms.
实验1
我在Amazon EMR主节点安装了Kafka CLI,可通过以下命令列出Amazon MSK中的Topic:
[hadoop@ip-xxx-xxx-xxx-xxx kafka_2.13-3.4.0]$ bin/kafka-topics.sh \ --bootstrap-server=b-1.myemr.xxx.c12.kafka.us-west-2.amazonaws.com:9098,b-2.myemr.xxx.c12.kafka.us-west-2.amazonaws.com:9098 \ --command-config=config/client.properties \ --list __amazon_msk_canary __consumer_offsets _schema_encoders _schemas hm.motor.avro
config/client.properties配置如下:
security.protocol=SASL_SSL sasl.mechanism=AWS_MSK_IAM sasl.jaas.config=software.amazon.msk.auth.iam.IAMLoginModule required; sasl.client.callback.handler.class=software.amazon.msk.auth.iam.IAMClientCallbackHandler
可见hm.motor.avro已存在,且通过相同IAM认证方式,EMR主节点可正常访问MSK的Topic。
另外,创建EMR集群时已设置实例配置文件为集群所有EC2实例分配角色,理论上EMR主节点和工作节点应拥有相同的MSK访问权限。
恳请各位提供解决思路,谢谢!
内容的提问来源于stack exchange,提问作者Hongbo Miao
相关产品推荐
相关产品推荐

