如何为Kinesis Consumer写入S3的文件添加前缀或后缀?
给Kinesis Stream写入S3的文件添加前缀/后缀的方法
嘿,很高兴你已经搞定了文件覆盖的问题!针对你想给S3中生成的文件添加前缀和后缀的需求,结合你当前用的Spark RDD API,这里有几个实用的方案:
方案1:保存后重命名文件(最直接易上手)
Spark的saveAsTextFile会默认生成类似part-00000的文件,我们可以先把文件存到临时目录,再用AWS SDK手动重命名成带前缀后缀的目标文件,最后清理临时文件。这种方法不需要修改Spark的输出逻辑,适合快速实现。
举个Scala代码示例(Java思路类似):
import software.amazon.awssdk.services.s3.S3Client import software.amazon.awssdk.services.s3.model.{CopyObjectRequest, DeleteObjectRequest} // 1. 定义临时存储路径和最终路径的基础结构 val baseOutputDir = "s3://your-bucket/output" val year = "2024" val month = "05" val day = "20" val hour = "14" val minute = "30" val tempDir = s"${baseOutputDir}/${year}/${month}/${day}/${hour}/${minute}/temp" val targetDir = s"${baseOutputDir}/${year}/${month}/${day}/${hour}/${minute}" // 2. 先把RDD保存到临时目录 rdd.coalesce(1).saveAsTextFile(tempDir) // 3. 初始化S3客户端(记得提前配置好AWS凭证和区域) val s3Client = S3Client.builder().build() val bucketName = "your-bucket" val tempPrefix = tempDir.replace(s"s3://${bucketName}/", "") // 4. 找到临时目录下的part文件(过滤掉_SUCCESS标记文件) val partFile = s3Client.listObjectsV2(req => req.bucket(bucketName).prefix(tempPrefix)) .contents() .stream() .filter(obj => obj.key().contains("part-") && obj.key().endsWith(".txt")) .findFirst() .orElseThrow(() => new RuntimeException("未找到临时目录下的part文件")) // 5. 定义带前缀后缀的目标文件名,比如前缀kinesis_data_,后缀_processed.txt val timestamp = System.currentTimeMillis() val targetKey = s"${targetDir}/kinesis_data_${timestamp}_processed.txt".replace(s"s3://${bucketName}/", "") // 6. 复制part文件到目标路径 val copyReq = CopyObjectRequest.builder() .sourceBucket(bucketName) .sourceKey(partFile.key()) .destinationBucket(bucketName) .destinationKey(targetKey) .build() s3Client.copyObject(copyReq) // 7. 清理临时文件和目录 s3Client.deleteObject(DeleteObjectRequest.builder().bucket(bucketName).key(partFile.key()).build()) s3Client.deleteObject(DeleteObjectRequest.builder().bucket(bucketName).key(tempPrefix).build())
方案2:自定义Hadoop OutputFormat(更优雅的底层方案)
Spark的saveAsTextFile底层依赖Hadoop的OutputFormat,我们可以自定义一个TextOutputFormat子类,修改文件名的生成逻辑,让它直接输出带前缀后缀的文件。这种方法不需要额外的重命名步骤,一步到位。
示例代码:
import org.apache.hadoop.fs.Path import org.apache.hadoop.io.{LongWritable, Text} import org.apache.hadoop.mapreduce.lib.output.TextOutputFormat import org.apache.hadoop.mapreduce.TaskAttemptContext // 自定义OutputFormat,重写文件名生成逻辑 class PrefixedTextOutputFormat extends TextOutputFormat[LongWritable, Text] { override def getDefaultWorkFile(context: TaskAttemptContext, extension: String): Path = { // 自定义前缀和后缀 val filePrefix = "kinesis_stream_" val fileSuffix = "_batch.txt" // 获取默认生成的文件名(比如part-00000) val defaultPath = super.getDefaultWorkFile(context, "") // 拼接成带前缀后缀的新路径 new Path(defaultPath.getParent, s"${filePrefix}${defaultPath.getName}${fileSuffix}") } } // 使用自定义OutputFormat保存RDD val outputPath = s"${baseOutputDir}/${year}/${month}/${day}/${hour}/${minute}" rdd.coalesce(1).saveAsNewAPIHadoopFile( outputPath, classOf[LongWritable], classOf[Text], classOf[PrefixedTextOutputFormat] )
这样生成的文件会直接是kinesis_stream_part-00000_batch.txt,同时保留原来的分区目录结构,也不会影响_SUCCESS文件的生成(如果不需要这个文件,可以后续手动删除)。
注意事项
- 如果你的RDD没有coalesce(1),会生成多个part文件,每个文件都会被加上前缀后缀,这点可以根据你的并行度需求调整。
- 不管用哪种方法,都要确保S3客户端有足够的权限(读、写、删除操作)。
- 自定义OutputFormat的方式更适合长期使用,减少额外的IO操作;而重命名方式更灵活,适合临时调整或需要动态生成前缀后缀的场景。
内容的提问来源于stack exchange,提问作者abhijeet
相关产品推荐
相关产品推荐

