【问题标题】:How can I read in a binary file from hdfs into a Spark dataframe?如何将二进制文件从 hdfs 读入 Spark 数据帧?
【发布时间】:2016-09-21 16:52:52
【问题描述】:

我正在尝试将一些代码从 pandas 移植到 (py)Spark。不幸的是,我已经在输入部分失败了,我想在其中读取二进制数据并将其放入 Spark Dataframe 中。

到目前为止,我正在使用来自 numpy 的fromfile

dt = np.dtype([('val1', '<i4'),('val2','<i4'),('val3','<i4'),('val4','f8')])
data = np.fromfile('binary_file.bin', dtype=dt)
data=data[1:]                                           #throw away header
df_bin = pd.DataFrame(data, columns=data.dtype.names)

但是对于 Spark,我找不到如何去做。到目前为止,我的解决方法是使用 csv-Files 而不是二进制文件,但这不是一个理想的解决方案。我知道我不应该将 numpy 的 fromfile 与 spark 一起使用。 如何读取已经加载到 hdfs 中的二进制文件?

我尝试了类似的东西

fileRDD=sc.parallelize(['hdfs:///user/bin_file1.bin','hdfs:///user/bin_file2.bin])
fileRDD.map(lambda x: ???)

但它给了我一个No such file or directory 错误。

我看过这个问题: spark in python: creating an rdd by loading binary data with numpy.fromfile 但这仅在我将文件存储在驱动程序节点的主目录中时才有效。

【问题讨论】:

    标签: python hadoop numpy apache-spark spark-dataframe


    【解决方案1】:

    因此,对于像我一样从 Spark 开始并偶然发现二进制文件的任何人。以下是我的解决方法:

    dt=np.dtype([('idx_metric','>i4'),('idx_resource','>i4'),('date','>i4'),
                 ('value','>f8'),('pollID','>i2')])
    schema=StructType([StructField('idx_metric',IntegerType(),False),
                       StructField('idx_resource',IntegerType(),False), 
                       StructField('date',IntegerType),False), 
                       StructField('value',DoubleType(),False), 
                       StructField('pollID',IntegerType(),False)])
    
    filenameRdd=sc.binaryFiles('hdfs://nameservice1:8020/user/*.binary')
    
    def read_array(rdd):
        #output=zlib.decompress((bytes(rdd[1])),15+32) # in case also zipped
        array=np.frombuffer(bytes(rdd[1])[20:],dtype=dt) # remove Header (20 bytes)
        array=array.newbyteorder().byteswap() # big Endian
        return array.tolist()
    
    unzipped=filenameRdd.flatMap(read_array)
    bin_df=sqlContext.createDataFrame(unzipped,schema)
    

    现在您可以在 Spark 中使用您的数据框做任何您想做的事情。

    【讨论】:

      【解决方案2】:

      编辑: 请查看此处提到的 sc.binaryFiles 的使用: https://stackoverflow.com/a/28753276/5088142


      尝试使用:

      hdfs://machine_host_name:8020/user/bin_file1.bin
      

      您是 core-site.xml

      fs.defaultFS 中的主机名

      【讨论】:

      • fs.defaultFS 说 nameservice1,但也有 hdfs://nameservice1:8020/user/bin_file1.bin 我仍然收到文件未找到错误。它可以与我在地图中放置的功能相关联吗? def read_bin: with open("myfile", "rb") as f: byte = f.read(1) while byte != "": byte = f.read(1)
      • 您在哪一行得到“找不到文件错误”?您打算如何使用“read_bin”功能? open 方法似乎不适用于 HDFS....
      • 错误在 read_bin 的第 2 行。没错,open 方法不喜欢 HDFS。我正在寻找类似于sc.textfile(filename).map(lambda line:line.split(',')).map(lambda x: (int(x[0],int(x[1]....)
      • 请查看此处提到的 sc.binaryFiles 的使用:stackoverflow.com/a/28753276/5088142
      【解决方案3】:

      从 Spark 3.0 开始,Spark 支持二进制文件数据源,即读取二进制文件并将每个文件转换为包含文件原始内容和元数据的单个记录。

      https://spark.apache.org/docs/latest/sql-data-sources-binaryFile.html

      【讨论】:

        【解决方案4】:

        我最近做了这样的事情:

        from struct import unpack_from
        
        # creates an RDD of binaryrecords for determinted record length
        binary_rdd = sc.binaryRecords("hdfs://" + file_name, record_length)
        
        # map()s each binary record to unpack() it
        unpacked_rdd = binary_rdd.map(lambda record: unpack_from(unpack_format, record))
        
        # registers a data frame with this schema; registerTempTable() it as table_name
        raw_df = sqlc.createDataFrame(unpacked_rdd, sparkSchema)
        raw_df.registerTempTable(table_name)
        

        其中 unpack_format 和 sparkSchema 必须“同步”。

        我有一个动态生成 unpack_format 和 sparkSchema 变量的脚本;两者同时。 (它是一个更大的代码库的一部分,所以为了便于阅读,这里不张贴)

        unpack_format 和 sparkSchema 可以定义如下,例如,

        from pyspark.sql.types import *
        
        unpack_format = '<'   # '<' means little-endian: https://docs.python.org/2/library/struct.html#byte-order-size-and-alignment
        sparkSchema = StructType()
        record_length = 0
        
        unpack_format += '35s'    # 35 bytes that represent a character string
        sparkSchema.add("FirstName", 'string', True)  # True = nullable
        record_length += 35
        
        unpack_format += 'H'    # 'H' = unsigned 2-byte integer
        sparkSchema.add("ZipCode", 'integer', True)
        record_length += 2
        
        # and so on for each field..
        

        【讨论】:

          猜你喜欢
          • 2015-09-09
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 2017-05-20
          • 2011-04-28
          • 1970-01-01
          • 2019-08-22
          • 2017-11-01
          相关资源
          最近更新 更多