Structured streaming默認(rèn)支持的sink類(lèi)型有File sink哪雕,F(xiàn)oreach sink,Console sink师枣,Memory sink血巍。
特別的說(shuō)明一下Foreach sink的用法(ps:以通過(guò)Foreach sink寫(xiě)入外部redis為例)。
lastEtlData.writeStream().foreach(new TestForeachWriter()).outputMode("update").start();
foreach方法的參數(shù)為ForeachWriter對(duì)象端逼,看下api說(shuō)明:
datasetOfString.writeStream().foreach(new ForeachWriter<String>() {
@Override
public boolean open(long partitionId, long version) {
// open connection
//此處用于打開(kāi)連接朗兵,以redis為例,此處從redis連接池中獲取連接
}
@Override
public void process(String value) {
// write string to connection
//此處用于數(shù)據(jù)寫(xiě)入redis顶滩。value為GenericRowWithSchema對(duì)象
}
@Override
public void close(Throwable errorOrNull) {
// close the connection
//此處用于關(guān)閉連接
}
});
看一下三個(gè)方法在的調(diào)用過(guò)程余掖,因?yàn)槭敲總€(gè)Partition的一批數(shù)據(jù)調(diào)用一次,還是需要關(guān)注下open礁鲁,close的頻率盐欺。一批數(shù)據(jù)open,close各一次仅醇。如下所示:
data.queryExecution.toRdd.foreachPartition { iter =>
if (writer.open(TaskContext.getPartitionId(), batchId)) {
try {
while (iter.hasNext) {
writer.process(encoder.fromRow(iter.next()))
}
} catch {
case e: Throwable =>
writer.close(e)
throw e
}
writer.close(null)
} else {
writer.close(null)
}
}
附寫(xiě)redis的foreachwriter實(shí)現(xiàn):
public static class TestForeachWriter extends ForeachWriter implements Serializable{
public static JedisPool jedisPool;
public Jedis jedis;
static {
JedisPoolConfig config = new JedisPoolConfig();
config.setMaxTotal(20);
config.setMaxIdle(5);
config.setMaxWaitMillis(1000);
config.setMinIdle(2);
config.setTestOnBorrow(false);
jedisPool = new JedisPool(config, "127.0.0.1", 6379);
}
public static synchronized Jedis getJedis() {
return jedisPool.getResource();
}
@Override
public boolean open(long partitionId, long version) {
jedis = getJedis();
return true;
}
@Override
public void process(Object value) {
GenericRowWithSchema genericRowWithSchema = (GenericRowWithSchema) value;
System.out.println(((GenericRowWithSchema) value).get(0).toString()+"-----------"+ ((GenericRowWithSchema) value).get(2).toString());
jedis.set(((GenericRowWithSchema) value).get(Integer.parseInt(genericRowWithSchema.schema().getFieldIndex("ID").get().toString())).toString(),((GenericRowWithSchema) value).get(Integer.parseInt(genericRowWithSchema.schema().getFieldIndex("COUNT(ADA1)").get().toString())).toString());
System.out.println("++++++++++"+((GenericRowWithSchema) value).get(Integer.parseInt(genericRowWithSchema.schema().getFieldIndex("ID").get().toString())).toString());
}
@Override
public void close(Throwable errorOrNull) {
jedis.close();
}
}
特別提醒:需要顯示實(shí)現(xiàn)Serializable接口冗美。