顯示具有 tcp 標籤的文章。 顯示所有文章
顯示具有 tcp 標籤的文章。 顯示所有文章

2019/01/28

Google TCP BBR Congestion Control

自 1980 年代網際網路崛起,到 1990 年代快速發展開始,就產生了 TCP 以及 UDP 的 IP 網路,1983年1月1日,ARPANet要求未來所有的網路傳輸都使用 IP 網路,統一了開放網路的規格。

網路通道就像是一條水管/水道,這條通道用來傳送資料,因為沒有即時調節流量的機制,TCP 透過接收端發送確認已經收到封包的 ACK,來判斷是否發送的速度太快,流量控制就是在控制資料發送端的發送速度。

BBR (Bottleneck Bandwidth and Round-trip propagation time) 演算法是 Google 在 2016 年提出的流量控制演算法,他透過有效頻寬偵測機制,並降低網路節點的 buffer 使用量,藉由減低重傳的效能消耗問題,提高頻寬使用率。

流量控制是在做什麼工作

把網路傳輸通道 TCP 想像為一條從山頂的水庫到目的地海洋的河流或水管,水庫有大量水資源,希望能運用水道以最短時間送到海洋,但是水庫並不知道這個水道有多寬,能以多少速度的水量放水。如果水放得太慢,無法享用水道的總流量,水放得太急,有可能會超過水道容量,而讓水漫出水道,造成淹水的情況。

流量控制,就是控制水庫放水速度的一個演算法,因為水道容水量的情況是隨時都在改變的,問題在於,並沒有即時回報的機制,可以隨時調節流量,另外,由於發送端貪婪的本性,他會希望盡可能完全佔用網路水道的所有容水量。

但總不能一直送一直送,也不管接收端有沒有收到資料吧。TCP 的機制是,在接收到收到封包時,會回覆一個 ACK 封包,確認已經收到了,透過 ACK 的偵測,就可以知道現在發送的速度,是不是已經超過了網路通道的容量,如果發生封包遺失的狀況時,發送端就知道要降低發送速度了。

雖然是開放網路,公平競爭,但實際上還是看誰最會搶佔網路資源,因此有些演算法強調搶佔的特性,會侵蝕掉使用其他演算法的 TCP 連線。

但從提供網路服務這一端來看,他希望所有來使用服務的使用者,可以公平地使用伺服器對外的網路通道,而且要盡可能減少因為 TCP 重傳的機制,造成無效的網路資源浪費,這時候,採用強調發送端公平使用網路的演算法,會比較有利。

因為網路通道本身不穩定的特性,如果中間有遇到無線網路時,這種情況會更嚴重,因此網路的路由器本身,通常都會加上 Buffer 的機制,能夠讓暫時無法發送的網路資料,存放在 Buffer 裡面,希望能藉此改善整體網路的效能。就像是水道中間,會加上一些滯洪池的機制,預防水量瞬間增加的問題。

然而這些 Buffer 卻因為 TCP 流量控制貪婪的本性,而被濫用,因為 TCP 希望自己能盡可能使用到網路的最大流量,所以也會盡可能將自己的封包,把路由器的 Buffer 塞滿,讓自己的速度更快,這也會造成 Bufferbloat 的問題。

網路上有許多網路節點,除了控制路由以外,有些節點還增加了 QoS(Quality of Service) 或是 Traffic Shaping 的機制,這也是基於剛剛提到的 Buffer 而提供的功能,因為有了 Buffer,路由器可以先將需要傳送的資料,放入不同等級的 Buffer 裡面,保留固定的傳送頻寬給具有高傳輸權限的網路封包。

BBR 演算法

BBR 演算法的細節,以這兩篇文章的說明比較清楚。

TCP BBR擁塞控制算法解析

Linux Kernel 4.9 中的 BBR 算法與之前的 TCP 擁塞控制相比有什麼優勢?

要不然就要看原始發表的論文

BBR: Congestion-Based Congestion Control Measuring bottleneck bandwidth and round-trip propagation time

BBR 解決問題的方案有兩點

  1. 因無法區分 congestion packet loss 及 error packet loss,BBR 考慮讓網路不產生 packet loss
  2. 因為把 buffer 塞滿可產生最大流量,但也會造成 RTT 降低,BBR 交替進行頻寬及 delay(RTT) 偵測

BBR 的四個狀態

  1. STARTUP 採用標準的 slow start 方式,指數增加發送速度,發現頻寬被佔滿時,就進入 DRAIN 的階段
  2. DRAIN 降低發送速度,將佔用的 buffer 排空
  3. PROBE_BW 改變發送速度進行頻寬偵測,在一個 RTT 內增加發送速度,如果 RTT 沒有改變,就降低發送速度,排空先前多送的封包,在六個 RTT 內使用這個發送速度
  4. PROBE_RTT 每經過 10s,如果沒有得到一個更低的延遲時間,就進入延遲偵測的階段,持續 200ms (或一個 RTT),這個階段固定發送 4 packets,偵測得到的最小 delay 時間作為最新的延遲時間

一些實測的結果

GCP採用新演算法TCP BBR 傳輸率將提高2700倍!

又一個 TCP BBR 的測試結果

spotify: Smoother Streaming with BBR

BBR 阻塞算法,真是黑科技

這些是使用了 BBR 以後的測試說明,全部都是正面,效能有改善的結果。

Google 在宣傳 BBR 時,都說明只要修改 Server 的部分,讓 Server 以 BBR 演算法運作,原因在於,流量控制演算法著眼的重點,是大量資料的發生源,會產生大量資料,發送出來的地方。只要修改 Server 的原因是他們是針對 Youtube 這樣提供串流服務的 Server 套用 BBR,換句話說,BBR 適用於 Sender Side。

因為串流影音的特性,就是需要一條長時間運作且傳輸量穩定的網路通道,這正好符合了 BBR 提供的流量結果,因為沒有使用到 Router 的 buffer,也沒有大量的重傳封包抵銷了網路的效能。

如果是類似 Hangouts Meeting 這樣的多人雙向影音的應用,因為客戶端 client side 如果沒有使用 BBR,就可能會產生不穩定的個人影音發生源,即使 Server Side 提供了 BBR,也無法形成一個有良好體驗的網路環境。

並不是所有人都認為 BBR 是有用的,以 令人躁動一時且令人不安的TCP BBR算法 這篇文章提出的論點來看,BBR 適合用在速度比較穩定的網路通道上,因為增速快,降速慢的特性,並不適用於忽快忽慢的網路。

我個人的想法是,只要決定了網路通道,固定了網路通道,那麼大部分的情況,網路是穩定的,該文章提出的問題,說明的並不恰當,如果通道上有些路徑的頻寬比較小,這會讓整條網路通道都因為這個最小頻寬的一段路而降速。因為流量控制只會根據接收端的 ACK 來調節,沒辦法知道中間經過每一個網路節點的速度。

會發生問題的地方,應該是網路忽快忽慢的情況,因為變化太大,演算法無法很快地調節到最佳的傳送速度,但這應該是所有演算法都會遇到的問題,BBR 改善的結果已經有顯著的效果了。

如何啟用 BBR

How to Deploy Google BBR on CentOS 7

開啟TCP BBR擁塞控制算法

因為 BBR 已經有在 linux kernel 4.9+ 的版本上實作,在各 linux distribution 的安裝方式,都是安裝新的 kernel repo,將 kernel 更新到 4.9+。

然後在 /etc/sysctl.conf 增加這兩行設定

net.core.default_qdisc = fq
net.ipv4.tcp_congestion_control = bbr
sysctl -p

用以下指令確認有沒有安裝成功

sysctl net.ipv4.tcp_available_congestion_control
sysctl net.ipv4.tcp_congestion_control
lsmod | grep bbr

BBR 可以終結流量控制的問題嗎?

這個 wiki 上的漫畫說明了現實的狀況,原本 BBR 的開發者,想要設計一個新的演算法,打敗既有12種 TCP Congestion Control 演算法,一統江湖,三年後,終於在 Linux 4.9 版 kernel 實現了 TCP BBR,但還是有某些缺陷,而現在變成了有 13 種 TCP Congestion Control 演算法。

https://upload.wikimedia.org/wikipedia/commons/3/34/Comic_strips_Linux_BBR.svg

References

TCP BBR

TCP擁塞控制

Faster Networking with TCP BBR

TCP BBR : Magic dust for network performance

Google最新tcp擁塞控制算法BBR解析

2014/04/14

erlang - Socket Programming

有兩個主要的函式庫:gen_tcp 與 gen_udp

從 web server 取得資料, tcp client

nano_get_url 可取得網頁的 html 內容

nano_get_url() ->
    nano_get_url("www.google.com").

nano_get_url(Host) ->
    %% 對 Host:80 開啟 TCP socket
    {ok,Socket} = gen_tcp:connect(Host,80,[binary, {packet, 0}]),
    %% 對此 Socket,傳送 GET / HTTP/1.0\r\n\r\n 字串的資料
    ok = gen_tcp:send(Socket, "GET / HTTP/1.0\r\n\r\n"),
    %% 等待接收server 回傳的資料
    receive_data(Socket, []).

receive_data(Socket, SoFar) ->
    receive
    %% 因為是用 binary 的方式開啟 socket,如果收到資料片段 Bin
    %% 就將它接到 SoFar 中暫存起來
    {tcp,Socket,Bin} ->
        receive_data(Socket, [Bin|SoFar]);
    %% 如果收到 tcp_closed 訊號,就表示 web server 送完資料,把 socket 斷線
    %% 就將資料 reverse,再將 binary list 連接成一整個 binary 資料
    {tcp_closed,Socket} ->
        list_to_binary(reverse(SoFar))
    end.

測試後,可取得 B 這一塊 binary 資料,然後可以用 io:format("~p~n", [B]) 將所有資料列印到畫面上,也可以用 string:tokens(binary_to_list(B), "\r\n") ,以 \r\n 切割資料,然後一行一行印出來。

1> B = socket_examples:nano_get_url().
<<"HTTP/1.0 302 Found\r\nLocation: http://www.google.com.tw/?gws_rd=cr&ei=W534Ur
qDKMLnkAWA-IHIDA\r\nCache-Control: private\r"...>>
2> io:format("~p~n", [B]).
<<"HTTP/1.0 302 Found\r\nLocation: http://www.google.com.tw/?gws_rd=cr&ei=W534Ur
qDKMLnkAWA-IHIDA\r\nCache-Control: private\r\nContent-Type: text/html; charset=U
TF-8\r\nSet-Cookie: PREF=ID=2d1f3f73d3fdb6d9:FF=0:TM=1392024923:LM=1392024923:S=
P-4nPCjpFPUCiXBs; expires=Wed, 10-Feb-2016 09:35:23 GMT; path=/; domain=.google.
...
3> string:tokens(binary_to_list(B), "\r\n").
["HTTP/1.0 302 Found",
 "Location: http://www.google.com.tw/?gws_rd=cr&ei=W534UrqDKMLnkAWA-IHIDA",
 "Cache-Control: private",
 "Content-Type: text/html; charset=UTF-8",
...

TCP Server: evaluate erlang expression, tcp server

需求是要有一個 tcp server, port 為 2345,他會等待 binary 訊息,裡面是一個 erlang expression,server運算後,將結果回傳給 client。

要想寫出任何一個 tcp 程式,必須先回答下列問題,因為 tcp socket 資料只是一個 bytes streaming,在傳輸期間,資料可以被打碎成任意長度的片段。

  1. 資料的格式,要如何知道 request 或 response 是由多少資料組成
  2. 在 request/response 內的資料,要如何編碼(marshaling)與解碼(de-marshaling)

{packet, N}

在 erlang 中,request/response可以用前置的 N(1或2或4) 個 bytes,來表示資料的長度,這也是 gen_tcp:connect 與 gentcp:listen 中 {packet, N} 參數的意義。

當我們利用 {packet, N} 開啟 socket 時,erlang driver 會自動將被打碎的資料片段,接合在一起。

term_to_binary 與 binary_to_term

erlang term 的編碼與解碼,可直接使用 term_to_binary 與 binary_to_term,這樣就不需要處理 http 或 xml 的文字編碼,不只速度快,傳送的資料也比較少。

程式

server 端的程式

start_nano_server() ->
    %% 開啟 tcp server port 2345,{packet, 4} 表示用 4 bytes 的資料長度 header
    %% gen_tcp:listen 會傳回 {ok, Socket} 或 {error, Why}
    {ok, Listen} = gen_tcp:listen(2345, [binary,     {packet, 4}, {reuseaddr, true}, {active, true}]),
    %% 將 Listen 綁定至 listen socket
    %% 在這裡,程式會暫停並等待 tcp client 連線
    {ok, Socket} = gen_tcp:accept(Listen),
    %% 當有 tcp client 連線後,就馬上關閉 Listen,這樣就不會再收到新的連線
    %% 且關閉後,不會影響既有的連線
    gen_tcp:close(Listen),
    %% 處理 socket 資料
    loop(Socket).

loop(Socket) ->
    receive
    {tcp, Socket, Bin} ->
        io:format("Server received binary = ~p~n",[Bin]),
        %% 將收到的資料 unmarshaling
        Str = binary_to_term(Bin),
        io:format("Server (unpacked)  ~p~n",[Str]),

        %% 估算 term
        Reply = lib_misc:string2value(Str),
        io:format("Server replying = ~p~n",[Reply]),

        %% 把結果 marshaling 之後,送進 socket
        gen_tcp:send(Socket, term_to_binary(Reply)),

        %% 等待此 tcp client 發送下一個 term,並處理
        loop(Socket);
    {tcp_closed, Socket} ->
        io:format("Server socket closed~n")
    end.

client 端的程式

nano_client_eval(Str) ->
    %% 開啟 socket,連接到 tcp port 2345
    {ok, Socket} = gen_tcp:connect("localhost", 2345,
            [binary, {packet, 4}]),
    %% 將 term 以 term_to_binary 編碼後,發送到 socket
    ok = gen_tcp:send(Socket, term_to_binary(Str)),

    %% 等待接收結果
    receive
    {tcp,Socket,Bin} ->
        io:format("Client received binary = ~p~n",[Bin]),
        %% 以 binary_to_term 解碼後,列印到畫面上
        Val = binary_to_term(Bin),
        io:format("Client result = ~p~n",[Val]),
        gen_tcp:close(Socket)
    end.

測試,啟動server時,畫面會停在這邊

1> socket_examples:start_nano_server().

啟動 client

1> socket_examples:nano_client_eval("list_to_tuple([2+3*4, 10+20])").

server 會收到資料,並估算結果

Server received binary = <<131,107,0,29,108,105,115,116,95,116,111,95,116,117,112,108,101,40,91,50,43,51,42,52,44,32,49,48,43,50,48,93,41>>
Server (unpacked)  "list_to_tuple([2+3*4, 10+20])"
Server replying = {14,30}
Server socket closed
ok

同時,client也會收到 server 回傳的結果

Client received binary = <<131,104,2,97,14,97,30>>
Client result = {14,30}
ok

改進 server

剛剛的 server 只會接受一個 client 連線,接下來嘗試修改server,讓它能接受多個連線。改進的方式有以下兩種

  1. 序列伺服器:一次接受一個連線
  2. 平行伺服器:同時接收多個連線
序列伺服器:一次接受一個連線

原本的程式為

start_nano_server() ->
    {ok, Listen} = gen_tcp:listen(2345, [binary, {packet, 4}, {reuseaddr, true}, {active, true}]),
    {ok, Socket} = gen_tcp:accept(Listen),
    gen_tcp:close(Listen),
    loop(Socket).

將程式改為以下的樣子,在 loop 完成後,繼續呼叫 seq_loop,讓它等候下一個連線。

start_seq_server() ->
    {ok, Listen} = gen_tcp:listen(2345, [binary, {packet, 4}, {reuseaddr, true}, {active, true}]),
    seq_loop(Listen).

seq_loop(Listen) ->
    {ok, Socket} = gen_tcp:accept(Listen),
    loop(Socket),
    seq_loop(Listen).

測試

1> socket_examples:start_seq_server().
Server received binary = <<131,107,0,29,108,105,115,116,95,116,111,95,116,117,
                           112,108,101,40,91,50,43,51,42,52,44,32,49,48,43,50,
                           48,93,41>>
Server (unpacked)  "list_to_tuple([2+3*4, 10+20])"
Server replying = {14,30}
Server socket closed
Server received binary = <<131,107,0,29,108,105,115,116,95,116,111,95,116,117,
                           112,108,101,40,91,50,43,51,42,52,44,32,50,48,43,50,
                           48,93,41>>
Server (unpacked)  "list_to_tuple([2+3*4, 20+20])"
Server replying = {14,40}
Server socket closed

client 的部份

1> socket_examples:nano_client_eval("list_to_tuple([2+3*4, 10+20])").
Client received binary = <<131,104,2,97,14,97,30>>
Client result = {14,30}
ok
2> socket_examples:nano_client_eval("list_to_tuple([2+3*4, 20+20])").
Client received binary = <<131,104,2,97,14,97,40>>
Client result = {14,40}
ok
平行伺服器:同時接收多個連線

在每一次 gen_tcp:accept 一產生新的連線後,就馬上 spawn 產生一個新的 process。

start_parallel_server() ->
    {ok, Listen} = gen_tcp:listen(2345, [binary, {packet, 4}, {reuseaddr, true}, {active, true}]),
    spawn(fun() -> par_connect(Listen) end).

par_connect(Listen) ->
    {ok, Socket} = gen_tcp:accept(Listen),
    spawn(fun() -> par_connect(Listen) end),
    loop(Socket).
註記
  1. 建立 socket 的行程(呼叫 gen_tcp:accept 或是 gen_tcp:connect )被稱為該 socket 的控制行程,來自 socket 的所有訊息都會送到控制行程中,如果控制行程死亡,socket 就會被關閉。可使用 gen_tcp:controlling_process(Socket, NewPid) 將控制行程換成 NewPid
  2. 平行伺服器可建立數千個連線,我們可能想限制同時連線的數量,可以使用一個 counter 記錄連線數。
  3. 接受連線後,最好明確設定 socket 選項:
     {ok, Socket} = gen_tcp:accept(Listen),
     inet:setopts(Socket, [{packet, 4}, binary, {nodelay,true}, {active, true}]),
     loop(Socket).
  4. 在 Erlang R11B-3,數個 erlang process 可對相同的 listen socket 呼叫 gen_tcp:accept/1,這可調整平行伺服器,用一些 pre-spawned processes pool,全部都等待 gen_tcp:accept/1。

控制 socket

erlang socket 可以三種模式開啟:active, active once, passive。active once 是建立主動 socket,收到訊息後,想接收下一個訊息,必須先重新啟用才行。

方式是在 gen_tcp:connect(Address, Port, Options) 或是 gen_tcp_listen(Port, Options) 的 Options 中使用 {active, true | false | once} 的設定。

主動與被動 socket 的差異:當 Socket 收到訊息時

  1. active socket 收到資料時,會送出 {tcp, Socket, Data} 給控制行程,控制行程無法控制這些訊息的流入。惡劣的客戶端就可能會送出數千個訊息給伺服器。
  2. passive socket 的控制行程必須呼叫 gen_tcp:recv(Socket, N),才能從 socket 接收 N bytes 的訊息,如果 N = 0,全部有效的 bytes data 都會被送出來。因此伺服器就能自行選擇何時呼叫 gen_tcp:recv ,這樣才能控制訊息流。

我們可用三種方式,撰寫 server 接收資料的迴圈

  1. active 訊息接收 - nonblocking
  2. passive 訊息接收 - blocking
  3. 混合訊息接收 - 部份 blocking
active 訊息接收 - nonblocking

process 無法控制進入 server 迴圈的訊息流,如果 client 端產生資料的速度比 server 消化還快,系統就會受到 message flooding,存放訊息的 mailbox 就有可能會 overflow,造成系統 crash。

因為無法阻塞 client,只有在我們確信能夠應付 client 的資料量時,才能使用 nonblocking server。

{ok, Listen} = gen_tcp:listen(Port, [..., {active, true}, ...]),
{ok, Socket} = gen_tcp:accept(Listen),
loop(Socket).

loop(Socket) ->
    receive
        {tcp, Socket, Data} ->
            ... do something with the data
        {tcp_closed, Socket} ->
            ...
    end.
passive 訊息接收 - blocking

server loop 只會在想要接收訊息時,才會呼叫 gen_tcp:recv。客戶端會被阻塞,直到伺服器呼叫 recv 為止。

注意:OS有做一些緩衝處理,即使尚未呼叫 recv,OS也會允許客戶端在被阻塞之前,可以送少量資料進來。

{ok, Listen} = gen_tcp:listen(Port, [..., {active, false}, ...]),
{ok, Socket} = gen_tcp:accept(Listen),
loop(Socket).

loop(Socket) ->
    case gen_tcp:recv(Socket, N) of
        {ok, B} ->
            ... do something with the data
            loop(Socket);
        {error, closed} ->
            ...
    end.
混合訊息接收 - 部份 blocking

只使用 passive mode 並不是最正確的作法,因為在 passive mode,只能等待一個 socket 的資料,對於必須等待多個 socket 資料的 server 來說,這是行不通的。

我們可以在開啟 Socket 時,使用 {active, once} ,這時候,Socket 只會對一個訊息主動。在控制行程收到一個訊息後,必須主動呼叫 inet:setopts 才能再次讓下一個訊息被接收,在此之前,系統會阻塞訊息。

使用 {active,once} ,使用者可實現 traffic shaping ,且避免 server 被過度積極湧入的訊息淹沒。

{ok, Listen} = gen_tcp:listen(Port, [..., {active, once}, ...]),
{ok, Socket} = gen_tcp:accept(Listen),
loop(Socket).

loop(Socket) ->
    receive
        {tcp, Socket, Data} ->
            ... do something with the data

            %% when you are ready to receive next message
            inet:setopts(Socket, [{active, once}]),
            loop(Socket);
        {tcp_closed, Socket} ->
            ...
    end.

連線來源

要知道 client 端的資訊,可以呼叫 inet:peername(Socket)

@spec inet:peername(Socket) -> {ok, {IP_Address, Port} | {error, Why}}

IP_Address 中
{N1, N2, N3, N4} 代表 IPv4
{K1, K2, K3, K4, K5, K6, K7, K8} 代表 IPv6

所有 Ni 與 Ki 都是介於 0 ~ 255 的整數

Socket 的錯誤處理

因為每一個 Socket 都有控制行程,當控制行程死亡,Socket 就會自動關閉。

error_test() ->
    %% 產生 server process
    spawn(fun() -> error_test_server() end),
    %% 先暫停兩秒,讓 server process 啟動
    lib_misc:sleep(2000),
    %% client 連線到 server
    {ok,Socket} = gen_tcp:connect("localhost",4321,[binary, {packet, 2}]),
    io:format("connected to:~p~n",[Socket]),
    %% 發送訊息 123
    gen_tcp:send(Socket, <<"123">>),
    receive
        %% 接收所有的 response 訊息
        Any ->
            io:format("Any=~p~n",[Any])
    end.

error_test_server() ->
    {ok, Listen} = gen_tcp:listen(4321, [binary,{packet,2}]),
    {ok, Socket} = gen_tcp:accept(Listen),
    error_test_server_loop(Socket).

error_test_server_loop(Socket) ->
    receive
        {tcp, Socket, Data} ->
            io:format("received:~p~n",[Data]),
            %% <<"123">> 會讓這一行當掉,server 的控制行程會 crash
            %% 因此會讓客戶端收到 {tcp_closed, Socket} 訊息
            atom_to_list(Data),
            error_test_server_loop(Socket)
    end.

測試

2> socket_examples:error_test().
connected to:#Port<0.594>
received:<<"123">>
Any={tcp_closed,#Port<0.594>}
ok
3>
=ERROR REPORT==== 11-Feb-2014::15:50:16 ===
Error in process <0.35.0> with exit value: {badarg,[{erlang,atom_to_list,[<<3 bytes>>],[]},{socket_examples,error_test_server_loop,1,[{file,"d:/projectcase/erlang/erlangotp/src/socket_examples.erl"},{line,117}]}]}

UDP server and client

UDP datagram 是不可靠的,順序可能被調換、可能會遺失、也可能會重複,且是 connectionless 的,客戶端不需要建立連線,就可以送訊息。

UDP 相當適合「大量客戶端,傳送小訊息到server」的應用情境。

UDP 比 TCP 簡單,因為 server 不需要管理與維護連線。

UDP server 一般的形式如下,server不會收到 socket 關閉的訊息:

server(Port) ->
    {ok, Socket} = gen_udp:open(Port, [binary]),
    loop(Socket).

loop(Socket) ->
    receive
        {udp, Socket, Host, Port, Bin} ->
            BinReply = ...,
            gen_udp:send(Socket, Host, Port, BinReply),
            loop(Socket)
    end.

client 必須要有 after timeout 機制,因為可能永遠等不到回應的訊息。

client(Request) ->
    {ok, Socket} = gen_udp:open(0, [binary]),
    ok = gen_udp:send(Socket, "localhost", 4000, Request),
    Value = receive
            {udp, Socket, _, _, Bin} ->
                {ok, Bin}
            after 2000 ->
                error
        end,
    gen_udp:close(Socket),
    Value.

UDP 階乘 server

範例程式

-module(udp_test).
-export([start_server/0, client/1]).

start_server() ->
    %% 產生 server process
    spawn(fun() -> server(4000) end).

%% The server           
server(Port) ->
    %% 開啟 udp server,接收 binary 資料
    {ok, Socket} = gen_udp:open(Port, [binary]),
    io:format("server opened socket:~p~n",[Socket]),
    loop(Socket).

loop(Socket) ->
    receive
        {udp, Socket, Host, Port, Bin} = Msg ->
            io:format("server received:~p~n",[Msg]),
            %% 將 binary 轉換為 erlang term
            N = binary_to_term(Bin),
            %% 運算階乘
            Fac = fac(N),
            %% 把結果回傳給 client
            gen_udp:send(Socket, Host, Port, term_to_binary(Fac)),
            %% 處理下一個訊息
            loop(Socket)
    end.

fac(0) -> 1;
fac(N) -> N * fac(N-1).

%% The client
client(N) ->
    {ok, Socket} = gen_udp:open(0, [binary]),
    io:format("client opened socket=~p~n",[Socket]),
    %% 以 term_to_binary 將 term 轉換為 binary
    ok = gen_udp:send(Socket, "localhost", 4000, 
                      term_to_binary(N)),
    Value = receive
                {udp, Socket, _, _, Bin} = Msg ->
                    io:format("client received:~p~n",[Msg]),
                    binary_to_term(Bin)
            after 2000 ->
                    0
            end,
    gen_udp:close(Socket),
    Value.

測試,server 的部份

1> udp_test:start_server().
server opened socket:#Port<0.516>
<0.33.0>
2> server received:{udp,#Port<0.516>,{127,0,0,1},54201,<<131,97,40>>}
2> server received:{udp,#Port<0.516>,{127,0,0,1},54202,<<131,97,20>>}
2> server received:{udp,#Port<0.516>,{127,0,0,1},60708,<<131,97,10>>}

client 的部份

1> udp_test:client(40).
client opened socket=#Port<0.516>
client received:{udp,#Port<0.516>,
                     {127,0,0,1},
                     4000,
                     <<131,110,20,0,0,0,0,0,64,37,5,255,100,222,15,8,126,242,
                       199,132,27,232,234,142>>}
815915283247897734345611269596115894272000000000
2> udp_test:client(20).
client opened socket=#Port<0.527>
client received:{udp,#Port<0.527>,
                     {127,0,0,1},
                     4000,
                     <<131,110,8,0,0,0,180,130,124,103,195,33>>}
2432902008176640000
3> udp_test:client(10).
client opened socket=#Port<0.528>
client received:{udp,#Port<0.528>,{127,0,0,1},4000,<<131,98,0,55,95,0>>}
3628800

補充說明

要注意,因為 UDP 是 connectionless 的協定,server無法藉由拒絕讀取資料而阻塞客戶端。

大型的 UDP packet 可能會被切成片段,以利在網路上傳輸,切割會發生在 router 所接受的 MTU(maximum transfer unit) 比 UDP packet 還小的時候 。一般會建議在 UDP 裡,一開始先使用較小的封包,然後慢慢增加資料量,並量測 throughput,如果 throughput 突然下降,就表示封包太大了。

UDP packet 有可能會被傳送兩次,所以在寫 RPC code 時要小心,不然可能就會執行兩次,傳回兩次。要避免發生這個問題時,可以加上 make_ref。

client(Request) ->
    {ok, Socket} -> gen_udp:open(0, [binary]),

    %% 產生唯一的識別參考
    Ref = make_ref(),
    B1 = term_to_binary(Ref, Request),
    ok = gen_udp:send(Socket, "localhost", 4000, B1),
    wait_for_ref(Socket, Ref).

wait_for_ref(Socket, Ref) ->
    receive
        {udp, Socket, _, _, Bin} ->
            case binary_to_term(Bin) of
                {Ref, Val} ->
                    Val;
                {_SomeoOtherRef, _} ->
                    wait_for_ref(Socket, Ref)
        end;
    after 1000 ->
        ...
    end.

廣播到多台機器

我們需要兩個 ports,一個用來送出廣播,一個用來傾聽廣播。

-module(broadcast).
-compile(export_all).

% 將IoList廣播到LAN的所有機器
% 負責廣播的 process 會開啟 port 5010
send(IoList) ->
    case inet:ifget("eth0", [broadaddr]) of
        {ok, [{broadaddr, Ip}]} ->
            {ok, S} =  gen_udp:open(5010, [{broadcast, true}]),
            gen_udp:send(S, Ip, 6000, IoList),
            gen_udp:close(S);
        _ ->
            io:format("Bad interface name, or\n"
                          "broadcasting not supported\n")
    end.

% 因為 windows 的 network interface 並沒有辦法直接取得 broadaddr
% 就直接把 broadcast address 寫在程式裡面
sendwindows(IoList) ->
    {ok, S} =  gen_udp:open(5010, [{broadcast, true}]),
    gen_udp:send(S, "192.168.1.255", 6000, IoList),
    gen_udp:close(S).

% 負責接收廣播的 process 會開啟 port 6000,並等待接收訊息
listen() ->
    {ok, _} = gen_udp:open(6000),
    loop().

loop() ->
    receive
        Any ->
            io:format("received:~p~n", [Any]),
            loop()
    end.

測試:server的部份

1> broadcast:listen().
received:{udp,#Port<0.516>,{192,168,1,57},5010,"{test"}

client

1> broadcast:sendwindows([123, "test"]).
ok

參考

Erlang and OTP in Action
Programming Erlang: Software for a Concurrent World