在大数据领域,实时写入(upsert)和快速OLAP查询一直是鱼和熊掌不能兼得,比如apache hudi,要事先决定好是倾向于快速写入还是快速OLAP查询,即Copy On Write Table vs. Merge On Read Table一旦选定好,就不能更改。databricks的delta-io也是类似的实现。而现实往往是希望在近乎实时upsert的同时,能快速的查询,至少是接近列存数据库的查询速度。正是这个需要,cloudera于2015年推出了apache kudu。这是一个支持schema的非K-V存储引擎,不依赖HDFS,内存是行存,flush到硬盘是列存,c++11实现,同时提供java/scala/python API,兼容Spark生态。弥补了大数据实时领域缺失的一环。

       本文实现一个demo,通过往kafka实时发送股票变化的event,模拟实时写入数据流。通过spark streaming消费kafka数据流,将数据实时写入kudu表,实现数据的入库。再用一个简单的jsp页面调用后端的spark SQL定时查询kudu表结果,展示一个近实时监控平台。由于牵涉多个apache开源项目,为了测试方便,都用docker或本地形式部署。最后可以通过browser查看结果,如下图所示。

realtime-monitor

实现步骤

生成kafka模拟数据

一个死循环,不停的往指定的kafka topic上写record,数据格式完全自己定义,用'\001'分割的key-value对。代码框架和主要数据结构如下:

// loop to generate data
       while(true) {
            while (openOrders.size() < 100) {
                Order newOrder = new Order();
                String newOrderSingleFIX = newOrder.newOrderSingleFIX();
                producer.send(new ProducerRecord<String, String>(kafkaTopic, newOrderSingleFIX));
                openOrders.add(newOrder);
            }

            for (Order order : openOrders) {
                if (!order.isComplete()) {
                    if (order.hasFurtherExecuted()) {
                        Thread.sleep(1);
                        String executionReportFIX = order.nextExecutionReportFIX();
                        producer.send(new ProducerRecord<String, String>(kafkaTopic, executionReportFIX));
                    }
                }
                else {
                    completedOrders.add(order);
                }
            }

            for (Order completedOrder : completedOrders) {
                openOrders.remove(completedOrder);
            }
            completedOrders.clear();

            //Thread.sleep(1);
        }
// Data
private static class Order {
        private String clordid;
        private String orderid;
        private int orderqty;
        private int leavesqty;
        private Symbol symbol;

        private final String pairDelimiter = "\001";
        private final String kvDelimiter = "=";

        private enum Symbol {
            AAPL, MSFT, ORCL, VMW, GOOG, AMZN, FB, TWTR
        }

        public Order() {
            clordid = UUID.randomUUID().toString();
            orderid = UUID.randomUUID().toString();
            orderqty = new Random().nextInt(10000);
            leavesqty = orderqty;
            symbol = Symbol.values()[new Random().nextInt(Symbol.values().length)];
        }
}

Spark Streaming消费kafka数据写往kudu

这部分为了方便,参考了spark-streaming里面和kafka有关的scala示例,简单修改一下代码,增加了parse从kafka拿到的event逻辑。目前只测试了spark跑在local的情况。

    var spark = SparkSession
      .builder()
      .appName("Kudu StockStreamer")
      .config("spark.some.config.option", "some-value")
      .getOrCreate()
    if (runLocal) {
      spark = SparkSession
        .builder()
        .appName("Kudu StockStreamer")
        .config("spark.some.config.option", "some-value")
        .config("spark.master", "local")
        .getOrCreate()
    }
    val sc = spark.sparkContext
    val sqlContext = spark.sqlContext
    val ssc = new StreamingContext(sc, Seconds(5))
    var kuduContext: KuduContext = new KuduContext(kuduMaster, sc)
    val broadcastSchema = sc.broadcast(schema)

    val topicSet = topics.split(",").toSet

    val kafkaParams = Map[String, Object](
      "bootstrap.servers" -> brokers,
      "key.deserializer" -> classOf[StringDeserializer],
      "value.deserializer" -> classOf[StringDeserializer],
      "group.id" -> s"kudu-stream-${System.currentTimeMillis()}"
    )
    val messages = KafkaUtils.createDirectStream[String, String](
      ssc, PreferConsistent,
      Subscribe[String, String](topicSet, kafkaParams))
    val parsed = messages.map(line => {
      parseFixEvent(line.value())
    })
    parsed.foreachRDD(rdd => {
      val df = sqlContext.createDataFrame(rdd,broadcastSchema.value)
      kuduContext.upsertRows(df,tableName)
    })

    sys.ShutdownHookThread {
      println("Gracefully stopping Spark Streaming Application")
      spark.close();
      ssc.stop(true, true)
      println("Application stopped")
    }

    // Start the computation
    ssc.checkpoint("./checkpoint")
    ssc.start()
    ssc.awaitTermination()

查询kudu结果的逻辑

为了调试方便,查询部分可以通过浏览器servlet查询,也可以通过命令行查询。核心逻辑是启动spark-SQL,查询kudu表。Query的参数决定查询最近多久的时间窗口,web.xml里面设置了600秒。

public class QueryTable {
    SparkSession _spark;
    String _sql_pat = "SELECT stocksymbol, max(orderqty) AS max_order," +
    "CAST((CAST(transacttime/10000 AS bigint)*10000)/1000 as timestamp) AS 10_s_time_window FROM " +
    "kafka_kudu WHERE transacttime > (CAST(unix_timestamp(to_utc_timestamp(now(),'PDT'))/10 AS bigint)*10 - %d)*1000 " +
            "AND transacttime < (CAST(unix_timestamp(to_utc_timestamp(now(),'PDT'))/10 AS bigint)*10 - 10)*1000 " + "" +
            "GROUP BY stocksymbol, 10_s_time_window ORDER BY stocksymbol, 10_s_time_window";
    public QueryTable(String kuduMasters, String kuduTable, String sparkMaster) {
        _spark = SparkSession
                .builder()
                .appName("Java Spark SQL basic example")
                .config("spark.some.config.option", "some-value")
                .config("spark.master", sparkMaster)
                .getOrCreate();
        Dataset<Row> df = _spark.read()
                .option("kudu.master", kuduMasters)
                .option("kudu.table", kuduTable)
                .option("kudu.scanLocality", "leader_only")
                .format("kudu").load();
        df.createOrReplaceTempView("kafka_kudu");
    }

    public void Stop() {
        _spark.close();
    }

    private String getSql(int seconds) {
        return String.format(_sql_pat, seconds);
    }

    public String Query(int seconds) {
        StringBuilder sb = new StringBuilder();
        Dataset<Row> result = _spark.sql(getSql(seconds));
        Dataset<String> formattedResult = result.map(
                (MapFunction<Row, String>)
                        row -> {
                            if (row.size() >= 3) {
                                return row.getString(0) + "," + row.getInt(1) + "," + row.getTimestamp(2);
                            } else {
                                return "";
                            }
                        },
                Encoders.STRING());
        sb.append("symbol,orderqty,timestamp").append(System.lineSeparator());
        Iterator<String> it = formattedResult.toLocalIterator();
        while (it.hasNext()) {
            sb.append(it.next());
            sb.append(System.lineSeparator());
        }
        return sb.toString();
    }
}

创建kudu表

这部分直接使用kudu的java API完成。唯一需要注意的是,这里根据时间创建了range partition,目的是为了提高OLAP扫描的速度。

def createFixTable(
    kuduMaster: String,
    tableName: String,
    numberOfHashPartitions: Int,
    numberOfDays: Int,
    tabletReplicas: Int): Unit = {

    val spark = SparkSession
      .builder()
      .appName("Spark-SQL kafka-kudu")
      .config("spark.some.config.option", "some-value")
      .getOrCreate()
    val kuduContext = new KuduContext(kuduMaster, spark.sparkContext)
    if(kuduContext.tableExists(tableName )) {
      System.out.println("Deleting existing table with same name.")
      kuduContext.deleteTable(tableName)
    }

    val options = new CreateTableOptions()
      .setRangePartitionColumns(ImmutableList.of("transacttime"))
      .addHashPartitions(ImmutableList.of("stocksymbol"),numberOfHashPartitions)
      .setNumReplicas(tabletReplicas)
    val today = new DateTime().withTimeAtStartOfDay() //today at midnight
    val dayInMillis = TimeUnit.MILLISECONDS.convert(1,TimeUnit.DAYS) //1 day in millis
    for (i <- 0 until numberOfDays) {
      val lbMillis = today.plusDays(i).getMillis
      val upMillis = lbMillis+dayInMillis-1
      val lowerBound = fixSchema.newPartialRow()
      lowerBound.addLong("transacttime",lbMillis)
      val upperBound = fixSchema.newPartialRow()
      upperBound.addLong("transacttime",upMillis)
      options.addRangePartition(lowerBound,upperBound)
    }
    kuduContext.createTable(tableName, schema, Seq("transacttime","stocksymbol","clordid"),options)

    System.out.println("Created new Kudu table " + tableName + " with " + numberOfHashPartitions + " hash partitions and " + numberOfDays + " date partitions. ")
    spark.close()
  }

测试须知

本地测试,请follow apache kudu的quick start,启动本地cluster。

参考https://gist.github.com/abacaphiliac/f0553548f9c577214d16290c2e751071,配置好kafka。

首先创建kudu表,然后开启consumer,再启动producer,最后mvn jetty:run启动jsp后端。

项目代码和测试步骤

https://github.com/ministat/KafkaMsgBroker

本文部分代码参考Building a Near RealTime Analytical Application with Kudu,并修正了Spark streaming API升级带来的不兼容。

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐