You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark-Scala中如何将S3/HDFS文件作为邮件附件发送?含DataFrame导出需求

Scala实现Spark DataFrame导出至S3/HDFS并作为邮件附件发送的最简方案

我来给你梳理下Scala里实现这个需求的最简方案,你之前遇到的问题核心是FileDataSource只能处理本地文件路径,没法直接读取S3/HDFS这类分布式文件系统里的文件,所以我们需要换个思路来获取文件流并构建邮件附件。

整体步骤拆解

  1. 将Spark DataFrame导出到S3/HDFS
  2. 通过Hadoop FileSystem API读取分布式文件的输入流
  3. 基于输入流构建邮件附件,用JavaMail发送邮件

1. 导出DataFrame到S3/HDFS

首先用Spark原生的write API就能完成导出,支持CSV、Parquet等多种格式,注意路径要使用分布式文件系统的协议前缀(比如s3a://或hdfs://):

import org.apache.spark.sql.SparkSession

// 初始化SparkSession
val spark = SparkSession.builder()
  .appName("DFExportAndEmail")
  // 如果是S3,这里可以提前配置访问密钥,或者在集群配置里设置
  .config("spark.hadoop.fs.s3a.access.key", "your-s3-access-key")
  .config("spark.hadoop.fs.s3a.secret.key", "your-s3-secret-key")
  .getOrCreate()

// 假设你的DataFrame已经准备好
val df = spark.read.table("your_source_table") // 替换成你的数据源

// 导出为CSV(带表头),也可以换成parquet等格式
val outputPath = "s3a://your-bucket/exported_data/output.csv" // 或hdfs://namenode:9000/exported_data/output.csv
df.write.mode("overwrite")
  .option("header", "true")
  .csv(outputPath)

注意:Spark导出分布式文件时会生成多个part-开头的分片文件,如果需要单个附件,建议先合并文件(比如用coalesce(1)强制生成单个文件,或者后续合并分片流)。


2. 读取S3/HDFS文件的输入流

用Hadoop的FileSystem API来访问分布式文件,获取文件输入流:

import org.apache.hadoop.fs.{FileSystem, Path}
import org.apache.hadoop.conf.Configuration

// 复用Spark的Hadoop配置(已经包含S3/HDFS的配置)
val hadoopConf = spark.sparkContext.hadoopConfiguration
val fs = FileSystem.get(new Path(outputPath).toUri, hadoopConf)

// 如果是多分片文件,先获取所有part文件的路径
val outputDir = new Path(outputPath)
val partFiles = fs.listStatus(outputDir)
  .filter(status => status.getPath.getName.startsWith("part-"))
  .map(_.getPath)

// 合并多个分片文件的输入流
import java.io.SequenceInputStream
import java.util.Collections
val combinedInputStream = new SequenceInputStream(Collections.enumeration(partFiles.map(fs.open(_)).asJava))

3. 构建邮件并发送附件

这里不用FileDataSource,而是把输入流转换成字节数组,用ByteArrayDataSource作为附件的数据源,再通过JavaMail API发送邮件:

import javax.mail._
import javax.mail.internet._
import javax.activation.ByteArrayDataSource

// 将输入流转换为字节数组
val fileBytes = Stream.continually(combinedInputStream.read())
  .takeWhile(_ != -1)
  .map(_.toByte)
  .toArray
combinedInputStream.close()

// 配置SMTP参数
val smtpProps = new java.util.Properties()
smtpProps.put("mail.smtp.host", "your-smtp-server.com") // 比如smtp.gmail.com
smtpProps.put("mail.smtp.port", "587")
smtpProps.put("mail.smtp.auth", "true")
smtpProps.put("mail.smtp.starttls.enable", "true")

// 创建邮件会话(带认证)
val session = Session.getInstance(smtpProps, new Authenticator() {
  override def getPasswordAuthentication(): PasswordAuthentication = {
    new PasswordAuthentication("your-email@example.com", "your-email-password")
  }
})

// 构建邮件内容
val message = new MimeMessage(session)
message.setFrom(new InternetAddress("your-email@example.com"))
message.setRecipients(Message.RecipientType.TO, InternetAddress.parse("recipient@example.com"))
message.setSubject("Spark DataFrame 导出附件")

// 多部分邮件内容(正文+附件)
val multipart = new MimeMultipart()

// 添加邮件正文
val bodyPart = new MimeBodyPart()
bodyPart.setText("这是Spark DataFrame导出的文件,请查收。")
multipart.addBodyPart(bodyPart)

// 添加附件
val attachmentPart = new MimeBodyPart()
// 根据文件类型设置MIME类型,比如CSV用"text/csv",Parquet用"application/octet-stream"
val dataSource = new ByteArrayDataSource(fileBytes, "text/csv")
attachmentPart.setDataHandler(new DataHandler(dataSource))
attachmentPart.setFileName("exported_data.csv") // 附件显示的文件名
multipart.addBodyPart(attachmentPart)

message.setContent(multipart)

// 发送邮件
Transport.send(message)

关键注意事项

  • 依赖配置:在你的build.sbt里需要添加以下依赖:
    libraryDependencies ++= Seq(
      "org.apache.spark" %% "spark-sql" % "3.3.0" % Provided,
      "javax.mail" % "mail" % "1.4.7",
      "org.apache.hadoop" % "hadoop-common" % "3.3.4" % Provided
    )
    
  • S3/HDFS权限:确保运行Spark的账号有对应的S3/HDFS读写权限,避免访问被拒绝。
  • 文件合并:如果不需要分片文件,导出时可以用df.coalesce(1).write...强制生成单个文件,但注意大数据量下coalesce(1)可能影响性能。

内容的提问来源于stack exchange,提问作者Raman Narasimhan

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.22 09:00:54