Flink1.32中使用Scala读取CSV转绑定Case Class的DataStream方法咨询
Flink 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解析即可避免拆分错误。
方案2:使用Flink 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
相关产品推荐
相关产品推荐

