Apache Kafka + Apache Kudu + Spark Streaming + Spark SQL实现大数据实时写入和实时监控
在大数据领域,实时写入(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查看结果,如下图所示。

实现步骤
生成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升级带来的不兼容。
更多推荐
所有评论(0)