Spark Streaming无法识别Azure Blob Storage中的application/octet-stream类型求助
解决Spark无法识别Azure Blob Storage中application/octet-stream类型文件的问题
嘿,我来帮你搞定这个问题!你遇到的核心问题是用错了InputFormat——TextInputFormat专门用来处理文本类文件,而你的Blob存储里是application/octet-stream格式的二进制文件,它自然没法解析这种二进制流内容。
给你几个针对性的解决方案:
方案1:用BinaryFileInputFormat读取二进制文件
如果你的文件是纯二进制格式,直接用BinaryFileInputFormat来读取,它能完整获取文件的二进制内容,代码示例如下:
import org.apache.hadoop.mapreduce.lib.input.BinaryFileInputFormat import org.apache.hadoop.io.{BytesWritable, NullWritable} // 读取二进制流文件,支持路径过滤 val dstream = ssc .fileStream[NullWritable, BytesWritable, BinaryFileInputFormat](filePath, (path: Path) => path.getName().endsWith(pathEndsWith), true) .map { case (_, bytesWritable) => // 将BytesWritable转换为字节数组,后续根据你的文件格式解析内容 val binaryContent = bytesWritable.getBytes // 这里可以添加自定义解析逻辑,比如转成字符串、解析为特定对象等,方便写入Cassandra binaryContent }
为什么这个方案可行?
BinaryFileInputFormat是Hadoop专门为二进制文件设计的输入格式,它返回的BytesWritable包含了文件的全部二进制数据,你可以根据实际文件的格式(比如自定义二进制协议、Protobuf、Avro等)来解析成Cassandra能接受的数据结构。
方案2:如果是序列文件,用SequenceFileInputFormat
如果你尝试的序列文件版本没写完,这里给你补全。如果Blob里的文件是Hadoop序列文件格式,就用SequenceFileInputFormat,代码示例:
import org.apache.hadoop.mapreduce.lib.input.SequenceFileInputFormat import org.apache.hadoop.io.LongWritable import org.apache.hadoop.io.Text // 注意这里的键值类型要和你的序列文件实际类型匹配,按需调整 val dstream = ssc .fileStream[LongWritable, Text, SequenceFileInputFormat[LongWritable, Text]](filePath) .map(_._2.toString)
额外注意点
- 确保你的Spark作业依赖了正确的Hadoop相关依赖包,避免出现类找不到的问题。
- 在写入Cassandra之前,一定要把解析后的二进制内容转换成对应的
case class或者Cassandra表的结构,这样Spark Cassandra Connector才能正确将数据写入到Cassandra中。
内容的提问来源于stack exchange,提问作者ivanpfalcao
相关产品推荐
相关产品推荐

