求助解决AWS Glue作业读取Glue Data Catalog中Kafka关联数据时的连接类型不支持错误
解决AWS Glue作业读取Kafka Catalog表时的"不支持KAFKA连接类型"错误
你遇到的问题根源很明确:Glue的from_catalog系列方法(包括create_data_frame.from_catalog和create_dynamic_frame.from_catalog)是为读取存储类数据源(比如S3、JDBC数据库)的Catalog表设计的,并不支持Kafka这类流式数据源的Catalog表。Data Catalog中注册的Kafka表仅保存了连接元数据,但无法直接通过from_catalog拉取数据。
正确解决方案:直接通过连接选项读取Kafka数据
你需要使用Glue专门针对外部数据源的create_dynamic_frame_from_options或create_data_frame_from_options方法,手动指定Kafka的连接参数。
示例1:创建Dynamic Frame读取Kafka
from awsglue.context import GlueContext from pyspark.context import SparkContext sc = SparkContext.getOrCreate() glueContext = GlueContext(sc) # 配置Kafka核心连接参数 kafka_options = { "bootstrap.servers": "你的Kafka集群地址:9092", "topicName": "目标Kafka主题名", "startingOffsets": "earliest", # 可选:latest/earliest或自定义偏移量 "inferSchema": "false", "security.protocol": "SSL" # 如果是加密集群需配置,根据实际环境调整 } # 创建Kafka动态帧 dynamic_frame = glueContext.create_dynamic_frame_from_options( connection_type="kafka", connection_options=kafka_options, transformation_ctx="kafka_dynamic_frame" )
示例2:创建Spark DataFrame读取Kafka
如果更习惯使用Spark DataFrame,可以用以下写法:
# 复用上面的kafka_options配置 kafka_df = glueContext.create_data_frame_from_options( connection_type="kafka", connection_options=kafka_options )
可选:从Data Catalog复用Kafka元数据参数
如果你想复用Data Catalog中已保存的Kafka连接配置,可以先通过Catalog API提取参数,再传入连接选项:
from pyspark.sql import SparkSession spark = SparkSession.builder.getOrCreate() # 从Catalog获取Kafka表的元数据 catalog_table = spark._jsparkSession.catalog().getTable("t-kafka-db.t-kafka") kafka_params = catalog_table.properties() # 整理Catalog中的参数到连接选项 kafka_options = { "bootstrap.servers": kafka_params.get("kafka.bootstrap.servers"), "topicName": kafka_params.get("kafka.topic.name"), "startingOffsets": "earliest", "inferSchema": "false" } # 后续创建Frame的步骤和示例1一致 dynamic_frame = glueContext.create_dynamic_frame_from_options( connection_type="kafka", connection_options=kafka_options, transformation_ctx="kafka_dynamic_frame" )
为什么之前的写法会报错?
当你调用from_catalog时,Glue会尝试根据Catalog表的连接类型初始化对应数据源的读取器,但目前from_catalog仅支持批处理类数据源(如S3、Redshift、JDBC等),流式数据源Kafka不在支持范围内,因此会抛出"We don't support this connection type: KAFKA"的错误。
内容的提问来源于stack exchange,提问作者Meruyert Mustafa
相关产品推荐
相关产品推荐

