基于RocksDB SST文件构建Spark DataFrame及性能优化问询
问题背景
我们将文档存储于RocksDB中,会把RocksDB的SST文件同步至S3,期望基于这些SST文件创建DataFrame并执行SQL查询,但未找到相关连接器。我们使用Spark 3.1.0搭配Scala 2.12,因将RocksDB转JSON再转Parquet的方式耗时且资源密集(需120个单核节点分钟),无法采用该方案。
进度更新
最终成功从AWS S3上的SST文件创建了DataFrame,但性能比基于JSON文件创建的DataFrame差4倍。
待优化功能
- 谓词下推(predicate pushdown):RocksDB的
_id字段是排序的,所有关联操作均基于该字段,实现谓词下推可大幅提升性能。
现有实现代码
Scala 部分
val ss = SparkSession.builder().getOrCreate(); val rows = //list of sst files from s3 var javaRdd = ss.sparkContext.parallelize(rows).toJavaRDD() javaRdd = javaRdd.flatMap(new KeyValueReader()); val scalaRdd: RDD[Row] = javaRdd.rdd val schemaJson = CloudUtils.getFileAsString(schemaUrl) val schema: StructType = DataType.fromJson(schemaJson).asInstanceOf[StructType] val df = ss.createDataFrame(scalaRdd, schema)
Java 部分
public class KeyValueReader implements FlatMapFunction<Row, Row> { @Override public Iterator<Row> call(Row row) throws Exception { return new KeyValueIterator(sstFile, schemaUrl); } }
public class KeyValueIterator implements Iterator<Row> { public boolean hasNext() { } public Row next() { //download sst file once. get the next row } }
内容的提问来源于stack exchange,提问作者chendu
相关产品推荐
相关产品推荐

