diff --git a/util/rdma/rdma_benchmark/client/client.go b/util/rdma/rdma_benchmark/client/client.go index 9c0a27005..91f946f57 100644 --- a/util/rdma/rdma_benchmark/client/client.go +++ b/util/rdma/rdma_benchmark/client/client.go @@ -1,6 +1,7 @@ package main import ( + "flag" "net" "rdma_test/common" "rdma_test/rdma" @@ -12,7 +13,7 @@ var Config = &rdma.RdmaEnvConfig{} func ReadBytes(conn net.Conn, buf []byte) error { offset := 0 - for offset < len(buf) { + for offset < len(buf) { n, err := conn.Read(buf[offset:]) if n == -1 || (err != nil) { println("ReadBytes failed, err: ", err) @@ -23,8 +24,19 @@ func ReadBytes(conn net.Conn, buf []byte) error { return nil } +func init() { + flag.StringVar(&Config.RdmaPort, "rdma-port", "9000", "rdma-port") + flag.IntVar(&Config.MemBlockNum, "memory-block-num", 0, "memory-block-num") + flag.IntVar(&Config.MemBlockSize, "memory-block-size", 0, "memory-block-size") + flag.IntVar(&Config.MemPoolLevel, "memory-pool-level", 0, "memory-pool-level") + flag.IntVar(&Config.ConnDataSize, "connect-data-size", 0, "connect-data-size") + flag.IntVar(&Config.WqDepth, "wq-depth", 0, "wq-depth") + flag.BoolVar(&Config.EnableRdmaLog, "enable-rdma-log", false, "enable-rdma-log") + flag.IntVar(&Config.WorkerNum, "worker-num", 0, "worker-num") + flag.StringVar(&Config.RdmaLogDir, "rdma-log-dir", "", "rdma-log-dir") +} -func testRdma() { +func testRdma() { if err := rdma.InitPool(Config); err != nil { println("init rdma pool failed") return @@ -34,14 +46,14 @@ func testRdma() { go func() { conn := &rdma.Connection{} - if err := conn.Dial(common.GParam.Ip,common.GParam.Port); err != nil { + if err := conn.Dial(common.GParam.Ip, common.GParam.Port); err != nil { println("client rdma conn dial failed") return } defer conn.Close() for !exit { beginTm := time.Now() - p := common.NewWritePacket(common.NormalExtentType, conn) + p := common.NewWritePacket(common.NormalExtentType, conn, uint32(common.GParam.IoSize)) p.Size = uint32(common.GParam.IoSize) if err := p.WriteToRDMAConn(conn); err != nil { @@ -61,13 +73,13 @@ func testRdma() { println(err.Error()) break } - err = rdma.ReleaseDataBuffer(conn, p.RdmaBuffer, uint32(105 + common.GParam.IoSize)) + err = rdma.ReleaseDataBuffer(conn, p.RdmaBuffer, uint32(105+common.GParam.IoSize)) if err != nil { println(err.Error()) break } - common.Stat().AddSumTime(common.GParam.IoSize, time.Now().UnixNano() / 1000 - beginTm.UnixNano() / 1000) + common.Stat().AddSumTime(common.GParam.IoSize, time.Now().UnixNano()/1000-beginTm.UnixNano()/1000) } }() } @@ -78,7 +90,7 @@ func testTcp() { for i := 0; i < common.GParam.IoDeep; i++ { go func() { - c, err := net.Dial("tcp", common.GParam.Ip + ":" + common.GParam.Port) + c, err := net.Dial("tcp", common.GParam.Ip+":"+common.GParam.Port) if err != nil { println(err.Error()) return @@ -89,7 +101,7 @@ func testTcp() { defer conn.Close() for !exit { beginTm := time.Now() - p := common.NewWritePacket(common.NormalExtentType, conn) + p := common.NewWritePacket(common.NormalExtentType, conn, uint32(common.GParam.IoSize)) p.Size = uint32(common.GParam.IoSize) if err := p.WriteToConn(conn); err != nil { @@ -103,14 +115,13 @@ func testTcp() { break } - common.Stat().AddSumTime(common.GParam.IoSize, time.Now().UnixNano() / 1000 - beginTm.UnixNano() / 1000) + common.Stat().AddSumTime(common.GParam.IoSize, time.Now().UnixNano()/1000-beginTm.UnixNano()/1000) } }() } } - -func main() { +func main() { common.ParseParam() if common.GParam.Protocol == "rdma" { go testRdma() @@ -123,4 +134,3 @@ func main() { common.Stat().Print() } } - diff --git a/util/rdma/rdma_benchmark/common/buffer_pool.go b/util/rdma/rdma_benchmark/common/buffer_pool.go index ea201f81c..fea7d92c2 100644 --- a/util/rdma/rdma_benchmark/common/buffer_pool.go +++ b/util/rdma/rdma_benchmark/common/buffer_pool.go @@ -2,10 +2,12 @@ package common import ( "fmt" + "math" "rdma_test/common/context" "rdma_test/common/rate" "sync" "sync/atomic" + "time" ) const ( @@ -74,6 +76,38 @@ func NewNormalBufferPool() *sync.Pool { } } +func NewBinaryPool() *BinaryPool { + pools := make([]*sync.Pool, 32) + for i := 0; i < 32; i++ { + pools[i] = &sync.Pool{ + New: func() interface{} { + buff := make([]byte, 1< 0 { - p.Arg = dataBuffer[65:65+p.ArgLen] + p.Arg = dataBuffer[65 : 65+p.ArgLen] } offset += 40 @@ -227,7 +227,7 @@ func (p *Packet) RecvRespFromRDMAConn(c *rdma.Connection, timeoutSec int) (err e return syscall.EBADMSG } size := p.Size - p.Data = dataBuffer[105:105+int(size)] + p.Data = dataBuffer[105 : 105+int(size)] return } @@ -290,7 +290,7 @@ func (p *Packet) ReadFromRDMAConnFromCli(conn *rdma.Connection, deadlineTime tim } offset += PacketHeaderSize + 8 //rdma if p.ArgLen > 0 { - p.Arg = dataBuffer[65:65+p.ArgLen] + p.Arg = dataBuffer[65 : 65+p.ArgLen] } offset += 40 if p.Size < 0 { @@ -299,11 +299,10 @@ func (p *Packet) ReadFromRDMAConnFromCli(conn *rdma.Connection, deadlineTime tim } size := p.Size - p.Data = dataBuffer[105:105+int(size)] + p.Data = dataBuffer[105 : 105+int(size)] return } - func ReadFull(c net.Conn, buf *[]byte, readSize int) (err error) { *buf = make([]byte, readSize) _, err = io.ReadFull(c, (*buf)[:readSize]) @@ -360,7 +359,6 @@ func (p *Packet) SendRespToRDMAConn(conn *rdma.Connection) (err error) { } }() - p.MarshalHeader(dataBuffer[0:PacketHeaderSize]) offset += PacketHeaderSize + 8 if p.ArgLen != 0 { @@ -376,4 +374,4 @@ func (p *Packet) SendRespToRDMAConn(conn *rdma.Connection) (err error) { return } return -} \ No newline at end of file +} diff --git a/util/rdma/rdma_benchmark/server/server.go b/util/rdma/rdma_benchmark/server/server.go index 87974d8d8..587b205e8 100644 --- a/util/rdma/rdma_benchmark/server/server.go +++ b/util/rdma/rdma_benchmark/server/server.go @@ -1,6 +1,7 @@ package main import ( + "flag" "net" "rdma_test/common" "rdma_test/rdma" @@ -12,7 +13,7 @@ var Config = &rdma.RdmaEnvConfig{} func ReadBytes(conn net.Conn, buf []byte) error { offset := 0 - for offset < len(buf) { + for offset < len(buf) { n, err := conn.Read(buf[offset:]) if n == -1 || (err != nil) { println("ReadBytes failed, err: ", err) @@ -23,7 +24,19 @@ func ReadBytes(conn net.Conn, buf []byte) error { return nil } -func testRdma() { +func init() { + flag.StringVar(&Config.RdmaPort, "rdma-port", "9000", "rdma-port") + flag.IntVar(&Config.MemBlockNum, "memory-block-num", 0, "memory-block-num") + flag.IntVar(&Config.MemBlockSize, "memory-block-size", 0, "memory-block-size") + flag.IntVar(&Config.MemPoolLevel, "memory-pool-level", 0, "memory-pool-level") + flag.IntVar(&Config.ConnDataSize, "connect-data-size", 0, "connect-data-size") + flag.IntVar(&Config.WqDepth, "wq-depth", 0, "wq-depth") + flag.BoolVar(&Config.EnableRdmaLog, "enable-rdma-log", false, "enable-rdma-log") + flag.IntVar(&Config.WorkerNum, "worker-num", 0, "worker-num") + flag.StringVar(&Config.RdmaLogDir, "rdma-log-dir", "", "rdma-log-dir") +} + +func testRdma() { if err := rdma.InitPool(Config); err != nil { println(err.Error()) return @@ -31,8 +44,8 @@ func testRdma() { server, _ := rdma.NewRdmaServer(common.GParam.Ip, common.GParam.Port) defer server.Close() - for { - conn,_ := server.Accept() + for { + conn, _ := server.Accept() go func() { for !exit { beginTm := time.Now() @@ -56,20 +69,19 @@ func testRdma() { break } - common.Stat().AddSumTime(common.GParam.IoSize, time.Now().UnixNano() / 1000 - beginTm.UnixNano() / 1000) + common.Stat().AddSumTime(common.GParam.IoSize, time.Now().UnixNano()/1000-beginTm.UnixNano()/1000) } conn.Close() }() } } - func testTcp() { common.InitBufferPool(0) - server, _ := net.Listen("tcp", common.GParam.Ip + ":" + common.GParam.Port) + server, _ := net.Listen("tcp", common.GParam.Ip+":"+common.GParam.Port) defer server.Close() - for { + for { c, _ := server.Accept() conn, _ := c.(*net.TCPConn) conn.SetKeepAlive(true) @@ -93,14 +105,13 @@ func testTcp() { break } - common.Stat().AddSumTime(common.GParam.IoSize, time.Now().UnixNano() / 1000 - beginTm.UnixNano() / 1000) + common.Stat().AddSumTime(common.GParam.IoSize, time.Now().UnixNano()/1000-beginTm.UnixNano()/1000) } }() } } - -func main() { +func main() { common.ParseParam() if common.GParam.Protocol == "rdma" { go testRdma()