Osiris是RabbitMQ使用的底層logs儲存機制,在 mix.exs 的 deps array 裡面加上
{:osiris, github: "rabbitmq/osiris", tag: "v1.13.1"}
來使用,hex.pm上面的同名專案是別的專案注意不要搞錯了。由於Osiris沒有公開API,要注意接下來講的用法未來可能被改變
準備一個資料夾給osiris,用Application env提供然後用 ensure_all_started 啟動osiris app
data_dir = ~c"/tmp/osiris"
File.mkdir_p!(data_dir)
:ok = Application.put_env(:osiris, :data_dir, data_dir)
{:ok, _apps} = Application.ensure_all_started(:osiris)
提供設定檔,單機上 leader_node 選目前的節點就可以了
config = %{
name: "my_stream",
epoch: 1,
leader_node: node(),
replica_nodes: []
}
{:ok, %{leader_pid: leader}} = :osiris.start_cluster(config)
寫入需要4個參數:osiris cluster的leader PID、writer ID、訊息序號、內容
writer_id = <<"writer-1">> :ok = :osiris.write(leader, writer_id, 1, <<"hello">>) :ok = :osiris.write(leader, writer_id, 2, <<"osiris">>)
成功寫入disk後會回一個 :osiris_written 訊息,所以寫入方要等這個消息傳過來,可以設定一個等待時限(這邊是5秒)
receive do
{:osiris_written, _name, ^writer_id, seqs} ->
IO.inspect(seqs, label: "written")
after
5_000 -> exit(:write_timeout)
end
讀取需要建立reader,會得到一個stream state,接下來就是用 read_chunk_parsed 解析出logs的內容
{:ok, log0} = :osiris.init_reader(leader, :first, {:my_reader, []})
{entries, log1} = :osiris_log.read_chunk_parsed(log0)
IO.inspect(entries, label: "chunk 1")
讀到 :end_of_stream 之後可以關閉logs
case :osiris_log.read_chunk_parsed(log1) do
{:end_of_stream, log2} ->
IO.puts("end of stream")
:osiris_log.close(log2)
{more, log2} ->
IO.inspect(more, label: "chunk 2")
:osiris_log.close(log2)
end
最後關閉cluster
:ok = :osiris.stop_cluster(config)