跳到主要內容

Java - NotificationBuffer With DelayQueue

Introduction

說到DelayQueue,我會想到拿來與ExecutorService一起使用,用以控制工作延遲多久被執行。為什麼我會把它與NotificationBuffer放在一起呢? 假設某一個Client有可能在短時間內送五筆相同資料過來,目前我們Server是會做五次例行工作;對於最後結果來說,其實我們只需要最後一次即可。除此之外,假如有多個Client,我希望能夠針對各別Client去計算這個時間區間。

Why DelayQueue?

從這兩點來說,我們需要一個buffer來儲存通知資料;然後這個buffer能夠根據各個Client第一次到達的時間,去決定何時讓Server取出並執行例行工作。以這樣的需求來說,DelayQueue是一個不錯的選擇。主要有以下幾個原因:

  • Producer只要把通知資料放進buffer,它就可以回傳了。
  • 繼承自BlockingQueue: 在buffer沒資料時,Consumer thread會處於block。這可以避免不必要的polling。
  • 使用PriorityQueue儲存資料: 這使得通知資料被加入buffer時,delay最短的那筆將會優先被take出去。這減少我們自己寫程式計算這些東西。
  • 根據delay時間去等待: 在每次資料被加入時,它會重新取得最短delay時間並做await。它使用必要的等待去取代不必要的polling檢查。

接下來我將說明要達成我目的的兩個重要物件:

  • NotifiedObject: 通知資料物件。
  • NotificationBuffer: 儲存通知資料物件的Buffer。

NotifiedObject

要實作自己的delayed物件,必須implements Dealyed介面,而它是extends Comparable介面,因此我們必須要實作compareTo與getDelay;除此之外,這個物件為了達到需求,會記載區別client的id與工作開始時間。可以參考底下程式碼:

import java.util.concurrent.Delayed;
import java.util.concurrent.TimeUnit;
 
public class NotifiedObject implements Delayed {
 
	private int oid;
	private long startTime;
 
	public NotifiedObject(int oid, int delayTime) {
		this.oid = oid;
		this.startTime = System.currentTimeMillis() + delayTime;
	}
 
	public int compareTo(Delayed o) {
	    return (int)(this.startTime - (((NotifiedObject)o).startTime));
	}
 
	@Override
	public long getDelay(TimeUnit unit) {
		long diff = startTime - System.currentTimeMillis();
		   return unit.convert(diff, TimeUnit.MILLISECONDS);
	}
 
	public int getHostObjId() {
		return oid;
	}
}

有幾個重點:

  • startTime在建立物件時,會根據目前時間與預期delay時間去算出來;delay時間就是你希望要緩衝的時間區間。
  • compareTo主要用途是給PriorityQueue決定先後順序用的,這可以讓DelayQueue取出最靠近現在時間的NotifiedObject。
  • getDelay是讓DelayQueue知道要等待多久才要將結果回應給client。

NotificationBuffer

為了讓特定client的多次通知合併成同一個,因此我做了一個extends DelayQueue的NotificationBuffer。在這個buffer中,idCache用以紀錄已經在NotifiedObject的oid,用以在add時略過不需要的通知資料;在take後,也會將oid從idCache中移除,以確保後來的通知能夠被處理:

import java.util.Set;
import java.util.concurrent.CopyOnWriteArraySet;
import java.util.concurrent.DelayQueue;
 
public static class NotificationBuffer extends DelayQueue<NotifiedObject> {
	private Set<Integer> idCache = new CopyOnWriteArraySet<Integer>();
 
	@Override
	public boolean add(NotifiedObject notifyObject) {
		if(idCache.contains(notifyObject.getHostObjId())) {
			return false;
		}
		try {
			return super.add(notifyObject);
		} finally {
			idCache.add(notifyObject.getHostObjId());
		}
	}
 
	@Override
	public NotifiedObject take() throws InterruptedException {
		NotifiedObject ret = super.take();
		idCache.remove(ret.getHostObjId());
		return ret;
	}
}

因為我僅對通知者有興趣,因此我可以將多出來的事件直接省略;假如你對於每次的通知內容都有興趣,你可以在這個物件中將多個NotifiedObject給合併成一個。

Unit Test

我的單元測試主要有兩個目標,一個是確認重複的通知有被濾掉,第二個則是確認buffer是可以被重複使用兩次以上的:

@Test(timeout=10*1000)
public void testDuplicatedNotifiedObject() throws Exception {
	NotificationBuffer buffer = new NotificationBuffer();
 
	NotifiedObject object1 = new NotifiedObject(111, 2000);
	NotifiedObject object2 = new NotifiedObject(111, 2000);
	NotifiedObject object3 = new NotifiedObject(333, 2001);
 
	buffer.add(object3);
	buffer.add(object2);
	buffer.add(object1);
 
	long before = System.currentTimeMillis();
	assertEquals(object2, buffer.take());
	assertEquals(2000, System.currentTimeMillis()-before, 10);
 
	assertEquals(object3, buffer.take());
	assertEquals(2000, System.currentTimeMillis()-before, 10);
 
	// re-add & take
	NotifiedObject object4 = new NotifiedObject(333, 2001);
	buffer.add(object4);
	before = System.currentTimeMillis();
	assertEquals(object4, buffer.take());
	assertEquals(2000, System.currentTimeMillis()-before, 10);
}
@Test(timeout=20*1000)
public void testTakeDelayTime() throws Exception {
	NotificationBuffer buffer = new NotificationBuffer();
 
	NotifiedObject object1 = new NotifiedObject(111, 2000);
	NotifiedObject object2 = new NotifiedObject(222, 4000);
	NotifiedObject object3 = new NotifiedObject(333, 6000);
 
	buffer.add(object3);
	buffer.add(object2);
	buffer.add(object1);
 
	long before = System.currentTimeMillis();
	assertEquals(object1, buffer.take());
	assertEquals(2000, System.currentTimeMillis()-before, 10);
 
	assertEquals(object2, buffer.take());
	assertEquals(4000, System.currentTimeMillis()-before, 10);
 
	assertEquals(object3, buffer.take());
	assertEquals(6000, System.currentTimeMillis()-before, 10);
}

後記

針對這樣的需求,Camel的Aggregator是可以達到部分的效果,可以參考這篇文章;我會說部分效果的原因,是因為Aggregator無法根據不同client到達的時間去計算delay,它是以第一個client到達的時間為準。至於要相依於框架還是要自己造輪子,就看各位的考量了。

Reference

留言

這個網誌中的熱門文章

解決RobotFramework從3.1.2升級到3.2.2之後,Choose File突然會整個Hand住的問題

考慮到自動測試環境的維護,我們很久以前就使用java去執行robot framework。前陣子開始處理從3.1.2升級到3.2.2的事情,主要先把明確的runtime語法錯誤與deprecate item處理好,這部分內容可以參考: link 。 直到最近才發現,透過SeleniumLibrary執行Choose File去上傳檔案的動作,會導致測試案例timeout。本篇文章主要分享心路歷程與解決方法,我也送了一條issue給robot framework: link 。 我的環境如下: RobotFramework: 3.2.2 Selenium: 3.141.0 SeleniumLibrary: 3.3.1 Remote Selenium Version: selenium-server-standalone-3.141.59 首先並非所有Choose File的動作都會hang住,有些測試案例是可以執行的,但是上傳一個作業系統ISO檔案一定會發生問題。後來我透過wireshark去比對新舊版本的上傳動作,因為我使用 Remote Selenium ,所以Selenium會先把檔案透過REST API發送到Remote Selenium Server上。從下圖我們可以發現,在3.2.2的最後一個TCP封包,比3.1.2大概少了500個bytes。 於是就開始了我trace code之路。包含SeleniumLibrary產生要送給Remote Selenium Server的request內容,還有HTTP Content-Length的計算,我都確認過沒有問題。 最後發現問題是出在socket API的使用上,就是下圖的這支code: 最後發現可能因為開始使用nio的方式送資料,但沒處理到尚未送完的資料內容,而導致發生問題。加一個loop去做計算就可以解決了。 最後我有把解法提供給robot framework官方,在他們出新的版本之前,我是將改完的_socket.py放在我們自己的Lib底下,好讓我們測試可以正常進行。(shutil.py應該也是為了解某個bug而產生的樣子..)

Show NIC selection when setting the network command with the device option

 Problem  在answer file中設定網卡名稱後,安裝時會停在以下畫面: 所使用的command參數如下: network --onboot = yes --bootproto =dhcp --ipv6 =auto --device =eth1 Diagnostic Result 這樣的參數,以前試驗過是可以安裝完成的。因此在發生這個問題後,我檢查了它的debug console: 從console得知,eth1可能是沒有連接網路線或者是網路太慢而導致的問題。後來和Ivy再三確認,有問題的是有接網路線的網卡,且問題是發生在activate階段: Solution 我想既然有retry應該就有次數或者timeout限制,因此發現在Anaconda的說明文件中( link ),有提到dhcptimeout這個boot參數。看了一些人的使用範例,應該是可以直接串在isolinux.cfg中,如下: default linux ksdevice = link ip =dhcp ks =cdrom: / ks.cfg dhcptimeout = 90 然而我在RHEL/CentOS 6.7與6.8試驗後都無效。 因此我就拿了顯示的錯誤字串,問問Google大師,想找一下Anaconda source code來看一下。最後找到別人根據Anaconda code修改的版本: link ,關鍵在於setupIfaceStruct函式中的setupIfaceStruct與readNetConfig: setupIfaceStruct: 會在dhcp時設定dhcptimeout。 readNetConfig: 在writeEnabledNetInfo將timeout寫入dhclient config中;在wait_for_iface_activation內會根據timeout做retry。 再來從log與code可以得知,它讀取的檔案是answer file而不是boot command line。因此我接下來的測試,就是在answer file的network command上加入dhcptimeout: network --onboot = yes --bootproto =dhcp --ipv6 =auto --device =eth1 --...

Robot Framework - Evaluate該怎麼用?

Evaluate該怎麼用? 前言 Builtin的RobotFramework Library提供了Evaluate Keyword。它所提供的功能是「執行Python描述句」。但實際上到底有什麼用途呢?原本我僅僅拿來將string轉為int的功用,經過一些查詢與試驗,我將心得整理給大家。 Builtin Builtin的function可以參考Library Doc for Evaluate。我以有使用過的function做說明。 數字轉換 Python提供了int、long、float與complex等function讓你可以將字串轉為數字,也可以透過它們做四則運算。首先以字串轉數字為例,我將8設於${num_str}中,再透過Evaluate+int轉為數字。這裡必須注意的是: 「int()中放變數必須以單引號'括起」。否則,假如你設定的數字為08,在轉換int時會出現Syntax Error。 ${num_str} | Set Variable | 8 ${num} | Evaluate | int('${num_str}') 其中int與long的第二個參數為base,這是根據你的input所決定: Comment | num = 9 ${num} | Evaluate | int('11', 8) Comment | num = 11 ${num} | Evaluate | int('11', 10) Comment | num = 17 ${num} | Evaluate | int('11', 16) 其它還有像bin、oct、hex,可以將整數轉為2、8、16進位。 運算 四則運算: 直接將運算子加上即可: ${num} | Evaluate | int('${hour}')*60 + int('${min}') 指數: 可以用pow。以下面兩個例子來說,第一個是2的3次方為8,第二個是2的3次方再mod 7為1。需注意的是: 「傳入值必須是數字不可為字串」。 ${num} | Evaluate | pow(2,3) ${num} | Evaluate | pow(2,3,7) 取最大最小值: 使用max/min,可以選擇丟一個array的方式...