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

Flink1.32中使用Scala读取CSV转绑定Case Class的DataStream方法咨询

方案1:基于现有String流做行转换(通用兼容所有版本)

你可以在现有读取到的String类型DataStream基础上,按CSV分隔符拆分后直接映射为你定义的Sales case class,实现代码如下:

import org.apache.flink.api.java.io.TextInputFormat
import org.apache.flink.api.scala.createTypeInformation
import org.apache.flink.core.fs.Path
import org.apache.flink.streaming.api.functions.source.FileProcessingMode
import org.apache.flink.streaming.api.scala.{DataStream, StreamExecutionEnvironment}


object AdvCsvRead {

  def main(args: Array[String]): Unit = {

    // 初始化执行环境
    val env = StreamExecutionEnvironment.getExecutionEnvironment

    val path = "src/main/resources/sales_orders.csv"
    val rawDs: DataStream[String] = env.readFile(new TextInputFormat(new Path(path)), path, FileProcessingMode.PROCESS_ONCE, 100)

    // 转换为Sales类型的DataStream,有表头的场景需增加过滤表头的逻辑
    val salesDs: DataStream[Sales] = rawDs
      .filter(line => !line.startsWith("ID,Customer")) // 替换为实际表头前缀,无表头可删除该行
      .map(line => {
        val fields = line.split(",") // 替换为实际CSV分隔符,复杂CSV建议用commons-csv库解析
        Sales(
          ID = fields(0).toInt,
          Customer = fields(1),
          Product = fields(2),
          Date = fields(3),
          Quantity = fields(4).toInt,
          Rate = fields(5).toDouble,
          Tags = fields(6)
        )
      })

    salesDs.print()
    env.execute("AdvCsvRead")
  }

  case class Sales (
                     ID: Integer,
                     Customer: String,
                     Product: String,
                     Date: String,
                     Quantity: Integer,
                     Rate: Double,
                     Tags: String
                   )

}

提示:如果CSV字段包含引号包裹的带逗号内容,不要直接用split拆分,引入org.apache.commons:commons-csv依赖做标准CSV解析即可避免拆分错误。

也可以直接使用Flink提供的CSV反序列化组件直接绑定字段读取,示例代码如下:

import org.apache.flink.api.scala._
import org.apache.flink.core.fs.Path
import org.apache.flink.streaming.api.functions.source.FileProcessingMode
import org.apache.flink.streaming.api.scala.{DataStream, StreamExecutionEnvironment}
import org.apache.flink.formats.csv.CsvRowDeserializationSchema
import org.apache.flink.types.Row

object AdvCsvRead {

  def main(args: Array[String]): Unit = {

    val env = StreamExecutionEnvironment.getExecutionEnvironment
    val path = "src/main/resources/sales_orders.csv"

    // 定义CSV反序列化规则,和Sales字段顺序、类型一一对应
    val deserializationSchema = CsvRowDeserializationSchema.builder()
      .setField("ID", classOf[Integer])
      .setField("Customer", classOf[String])
      .setField("Product", classOf[String])
      .setField("Date", classOf[String])
      .setField("Quantity", classOf[Integer])
      .setField("Rate", classOf[Double])
      .setField("Tags", classOf[String])
      .setIgnoreParseErrors(true) // 解析错误时跳过该行,不需要可关闭
      .setSkipFirstLineAsHeader(true) // 跳过表头,无表头可关闭
      .build()

    // 直接读取映射为Sales类型DataStream
    val salesDs: DataStream[Sales] = env.readFile(
      inputFormat = deserializationSchema,
      filePath = path,
      watchType = FileProcessingMode.PROCESS_ONCE,
      interval = 100
    ).map(row => Sales(
      ID = row.getFieldAs[Integer]("ID"),
      Customer = row.getFieldAs[String]("Customer"),
      Product = row.getFieldAs[String]("Product"),
      Date = row.getFieldAs[String]("Date"),
      Quantity = row.getFieldAs[Integer]("Quantity"),
      Rate = row.getFieldAs[Double]("Rate"),
      Tags = row.getFieldAs[String]("Tags")
    ))

    salesDs.print()
    env.execute("AdvCsvRead")
  }

  case class Sales (
                     ID: Integer,
                     Customer: String,
                     Product: String,
                     Date: String,
                     Quantity: Integer,
                     Rate: Double,
                     Tags: String
                   )

}
  • 方案1代码更简洁,适合小文件、CSV格式规则简单的场景
  • 方案2处理性能更高,内置标准CSV解析逻辑,适合格式复杂、文件体量较大的场景

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 01:27:03