使用AWS Glue(Spark)写入Kinesis时遇数据源未找到错误求助
AWS Glue 4.0(Scala)批量写入Kinesis报错:Failed to find data source: kinesis
问题场景
从S3/静态数据库表批量提取数据,经转换得到DynamicFrame后尝试写入Kinesis流,使用glueContext.getSinkWithFormat方法时触发错误:
Error writing to Kinesis: Failed to find data source: kinesis. Please find packages at https://spark.apache.org/third-party-projects.html
环境信息:Glue v4.0、Scala 2.12.19、Spark 3.3
原代码:
// 构建源并完成转换,最终得到DynamicFrame:myDynamicFrame val kinesis = glueContext.getSinkWithFormat( conectionType = "kinesis", options = JsonOptions( Map( "streamArn" -> "arn:aws:kinesis:xxxxxxxxxxx/sink-stream", "startingPosition" -> "TRIM_HORIZON", "inferSchema" -> "true" ) ) ) kinesis.writeDynamicFrame(myDynamicFrame)
问题原因
- 参数拼写错误:
conectionType应为connectionType(缺少字母'n') - 遗漏必填参数:
getSinkWithFormat方法需要指定输出format(如json),原代码未传入 - 无效参数:
startingPosition是Kinesis作为数据源时的配置,作为写入目标时无需设置
修复方案
方案1:修正getSinkWithFormat参数
调整代码中的拼写错误、补全必填参数并移除无效配置:
// 构建源并完成转换,最终得到DynamicFrame:myDynamicFrame val kinesisSink = glueContext.getSinkWithFormat( connectionType = "kinesis", // 修正拼写 options = JsonOptions( Map( "streamArn" -> "arn:aws:kinesis:xxxxxxxxxxx/sink-stream", "inferSchema" -> "true", "partitionKey" -> "your-partition-key-column" // 可选:指定分区键字段,无则可省略 ) ), format = "json" // 必须指定输出格式,如json、csv等 ) kinesisSink.writeDynamicFrame(myDynamicFrame)
方案2:转换为DataFrame直接写入
若方案1仍有问题,可将DynamicFrame转换为Spark DataFrame,使用Spark原生Kinesis写入API:
// 将DynamicFrame转为DataFrame val targetDF = myDynamicFrame.toDF() // 批量写入Kinesis targetDF.write .format("kinesis") .option("streamArn", "arn:aws:kinesis:xxxxxxxxxxx/sink-stream") .option("region", "your-aws-region") // 如us-east-1 .option("partitionKey", "your-partition-key-column") // 可选 .mode("append") .save()
额外注意事项
- 确保Glue作业关联的IAM角色拥有
kinesis:PutRecord和kinesis:PutRecords权限 - 若使用streamArn而非streamName,无需额外指定region(Arn已包含区域信息)
内容的提问来源于stack exchange,提问作者Richard Chester
相关产品推荐
相关产品推荐

