Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
milvus
milvus
提交
1601a61b
M
milvus
项目概览
milvus
/
milvus
9 个月 前同步成功
通知
260
Star
22476
Fork
2472
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
0
列表
看板
标记
里程碑
合并请求
0
DevOps
流水线
流水线任务
计划
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
M
milvus
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
0
Issue
0
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
DevOps
DevOps
流水线
流水线任务
计划
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
流水线任务
提交
Issue看板
前往新版Gitcode,体验更适合开发者的 AI 搜索 >>
未验证
提交
1601a61b
编写于
1月 14, 2022
作者:
Z
zhenshan.cao
提交者:
GitHub
1月 14, 2022
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
Move logic of checking available port to service (#15222)
Signed-off-by:
N
zhenshan.cao
<
zhenshan.cao@zilliz.com
>
上级
8873be6a
变更
8
隐藏空白更改
内联
并排
Showing
8 changed file
with
48 addition
and
44 deletion
+48
-44
internal/distributed/datanode/service.go
internal/distributed/datanode/service.go
+4
-1
internal/distributed/indexnode/service.go
internal/distributed/indexnode/service.go
+4
-1
internal/distributed/proxy/service.go
internal/distributed/proxy/service.go
+4
-5
internal/distributed/querynode/service.go
internal/distributed/querynode/service.go
+5
-0
internal/util/funcutil/func.go
internal/util/funcutil/func.go
+23
-0
internal/util/funcutil/func_test.go
internal/util/funcutil/func_test.go
+8
-0
internal/util/paramtable/global_param.go
internal/util/paramtable/global_param.go
+0
-29
internal/util/paramtable/global_param_test.go
internal/util/paramtable/global_param_test.go
+0
-8
未找到文件。
internal/distributed/datanode/service.go
浏览文件 @
1601a61b
...
@@ -217,7 +217,10 @@ func (s *Server) Stop() error {
...
@@ -217,7 +217,10 @@ func (s *Server) Stop() error {
func
(
s
*
Server
)
init
()
error
{
func
(
s
*
Server
)
init
()
error
{
ctx
:=
context
.
Background
()
ctx
:=
context
.
Background
()
Params
.
InitOnce
(
typeutil
.
DataNodeRole
)
Params
.
InitOnce
(
typeutil
.
DataNodeRole
)
if
!
funcutil
.
CheckPortAvailable
(
Params
.
Port
)
{
Params
.
Port
=
funcutil
.
GetAvailablePort
()
log
.
Warn
(
"DataNode get available port when init"
,
zap
.
Int
(
"Port"
,
Params
.
Port
))
}
dn
.
Params
.
InitOnce
()
dn
.
Params
.
InitOnce
()
dn
.
Params
.
DataNodeCfg
.
Port
=
Params
.
Port
dn
.
Params
.
DataNodeCfg
.
Port
=
Params
.
Port
dn
.
Params
.
DataNodeCfg
.
IP
=
Params
.
IP
dn
.
Params
.
DataNodeCfg
.
IP
=
Params
.
IP
...
...
internal/distributed/indexnode/service.go
浏览文件 @
1601a61b
...
@@ -119,7 +119,10 @@ func (s *Server) startGrpcLoop(grpcPort int) {
...
@@ -119,7 +119,10 @@ func (s *Server) startGrpcLoop(grpcPort int) {
func
(
s
*
Server
)
init
()
error
{
func
(
s
*
Server
)
init
()
error
{
var
err
error
var
err
error
Params
.
InitOnce
(
typeutil
.
IndexNodeRole
)
Params
.
InitOnce
(
typeutil
.
IndexNodeRole
)
if
!
funcutil
.
CheckPortAvailable
(
Params
.
Port
)
{
Params
.
Port
=
funcutil
.
GetAvailablePort
()
log
.
Warn
(
"IndexNode get available port when init"
,
zap
.
Int
(
"Port"
,
Params
.
Port
))
}
indexnode
.
Params
.
InitOnce
()
indexnode
.
Params
.
InitOnce
()
indexnode
.
Params
.
IndexNodeCfg
.
Port
=
Params
.
Port
indexnode
.
Params
.
IndexNodeCfg
.
Port
=
Params
.
Port
indexnode
.
Params
.
IndexNodeCfg
.
IP
=
Params
.
IP
indexnode
.
Params
.
IndexNodeCfg
.
IP
=
Params
.
IP
...
...
internal/distributed/proxy/service.go
浏览文件 @
1601a61b
...
@@ -158,12 +158,11 @@ func (s *Server) Run() error {
...
@@ -158,12 +158,11 @@ func (s *Server) Run() error {
func
(
s
*
Server
)
init
()
error
{
func
(
s
*
Server
)
init
()
error
{
Params
.
InitOnce
(
typeutil
.
ProxyRole
)
Params
.
InitOnce
(
typeutil
.
ProxyRole
)
log
.
Debug
(
"
init Proxy
service's parameter table done"
)
log
.
Debug
(
"
Proxy init
service's parameter table done"
)
proxy
.
Params
.
ProxyCfg
.
NetworkAddress
=
Params
.
GetAddress
()
if
!
funcutil
.
CheckPortAvailable
(
Params
.
Port
)
{
if
!
paramtable
.
CheckPortAvailable
(
Params
.
Port
)
{
Params
.
Port
=
funcutil
.
GetAvailablePort
()
// as the entry of Milvus, we'd better not to use another port
log
.
Warn
(
"Proxy get available port when init"
,
zap
.
Int
(
"Port"
,
Params
.
Port
))
return
fmt
.
Errorf
(
"port %d already in use"
,
Params
.
Port
)
}
}
proxy
.
Params
.
InitOnce
()
proxy
.
Params
.
InitOnce
()
...
...
internal/distributed/querynode/service.go
浏览文件 @
1601a61b
...
@@ -87,6 +87,11 @@ func NewServer(ctx context.Context, factory msgstream.Factory) (*Server, error)
...
@@ -87,6 +87,11 @@ func NewServer(ctx context.Context, factory msgstream.Factory) (*Server, error)
func
(
s
*
Server
)
init
()
error
{
func
(
s
*
Server
)
init
()
error
{
Params
.
InitOnce
(
typeutil
.
QueryNodeRole
)
Params
.
InitOnce
(
typeutil
.
QueryNodeRole
)
if
!
funcutil
.
CheckPortAvailable
(
Params
.
Port
)
{
Params
.
Port
=
funcutil
.
GetAvailablePort
()
log
.
Warn
(
"QueryNode get available port when init"
,
zap
.
Int
(
"Port"
,
Params
.
Port
))
}
qn
.
Params
.
InitOnce
()
qn
.
Params
.
InitOnce
()
qn
.
Params
.
QueryNodeCfg
.
QueryNodeIP
=
Params
.
IP
qn
.
Params
.
QueryNodeCfg
.
QueryNodeIP
=
Params
.
IP
qn
.
Params
.
QueryNodeCfg
.
QueryNodePort
=
int64
(
Params
.
Port
)
qn
.
Params
.
QueryNodeCfg
.
QueryNodePort
=
int64
(
Params
.
Port
)
...
...
internal/util/funcutil/func.go
浏览文件 @
1601a61b
...
@@ -22,7 +22,9 @@ import (
...
@@ -22,7 +22,9 @@ import (
"errors"
"errors"
"fmt"
"fmt"
"io/ioutil"
"io/ioutil"
"net"
"net/http"
"net/http"
"strconv"
"time"
"time"
"github.com/go-basic/ipv4"
"github.com/go-basic/ipv4"
...
@@ -172,3 +174,24 @@ func GetAttrByKeyFromRepeatedKV(key string, kvs []*commonpb.KeyValuePair) (strin
...
@@ -172,3 +174,24 @@ func GetAttrByKeyFromRepeatedKV(key string, kvs []*commonpb.KeyValuePair) (strin
func
CheckCtxValid
(
ctx
context
.
Context
)
bool
{
func
CheckCtxValid
(
ctx
context
.
Context
)
bool
{
return
ctx
.
Err
()
!=
context
.
DeadlineExceeded
&&
ctx
.
Err
()
!=
context
.
Canceled
return
ctx
.
Err
()
!=
context
.
DeadlineExceeded
&&
ctx
.
Err
()
!=
context
.
Canceled
}
}
// CheckPortAvailable check if a port is available to be listened on
func
CheckPortAvailable
(
port
int
)
bool
{
addr
:=
":"
+
strconv
.
Itoa
(
port
)
listener
,
err
:=
net
.
Listen
(
"tcp"
,
addr
)
if
listener
!=
nil
{
listener
.
Close
()
}
return
err
==
nil
}
// GetAvailablePort return an available port that can be listened on
func
GetAvailablePort
()
int
{
listener
,
err
:=
net
.
Listen
(
"tcp"
,
":0"
)
if
err
!=
nil
{
panic
(
err
)
}
defer
listener
.
Close
()
return
listener
.
Addr
()
.
(
*
net
.
TCPAddr
)
.
Port
}
internal/util/funcutil/func_test.go
浏览文件 @
1601a61b
...
@@ -282,3 +282,11 @@ func TestCheckCtxValid(t *testing.T) {
...
@@ -282,3 +282,11 @@ func TestCheckCtxValid(t *testing.T) {
time
.
Sleep
(
timeout
+
deltaTime
)
time
.
Sleep
(
timeout
+
deltaTime
)
assert
.
False
(
t
,
CheckCtxValid
(
ctx3
))
assert
.
False
(
t
,
CheckCtxValid
(
ctx3
))
}
}
func
TestCheckPortAvailable
(
t
*
testing
.
T
)
{
num
:=
10
for
i
:=
0
;
i
<
num
;
i
++
{
port
:=
GetAvailablePort
()
assert
.
Equal
(
t
,
CheckPortAvailable
(
port
),
true
)
}
}
internal/util/paramtable/global_param.go
浏览文件 @
1601a61b
...
@@ -13,7 +13,6 @@ package paramtable
...
@@ -13,7 +13,6 @@ package paramtable
import
(
import
(
"math"
"math"
"net"
"os"
"os"
"path"
"path"
"strconv"
"strconv"
...
@@ -23,7 +22,6 @@ import (
...
@@ -23,7 +22,6 @@ import (
"github.com/go-basic/ipv4"
"github.com/go-basic/ipv4"
"github.com/milvus-io/milvus/internal/log"
"github.com/milvus-io/milvus/internal/log"
"github.com/milvus-io/milvus/internal/util/typeutil"
"go.uber.org/zap"
"go.uber.org/zap"
)
)
...
@@ -1472,12 +1470,6 @@ func (p *grpcConfig) LoadFromArgs() {
...
@@ -1472,12 +1470,6 @@ func (p *grpcConfig) LoadFromArgs() {
func
(
p
*
grpcConfig
)
initPort
()
{
func
(
p
*
grpcConfig
)
initPort
()
{
p
.
Port
=
p
.
ParseInt
(
p
.
Domain
+
".port"
)
p
.
Port
=
p
.
ParseInt
(
p
.
Domain
+
".port"
)
if
p
.
Domain
==
typeutil
.
ProxyRole
||
p
.
Domain
==
typeutil
.
DataNodeRole
||
p
.
Domain
==
typeutil
.
IndexNodeRole
||
p
.
Domain
==
typeutil
.
QueryNodeRole
{
if
!
CheckPortAvailable
(
p
.
Port
)
{
p
.
Port
=
GetAvailablePort
()
log
.
Warn
(
"get available port when init"
,
zap
.
String
(
"Domain"
,
p
.
Domain
),
zap
.
Int
(
"Port"
,
p
.
Port
))
}
}
}
}
// GetAddress return grpc address
// GetAddress return grpc address
...
@@ -1620,24 +1612,3 @@ func (p *GrpcClientConfig) initClientMaxRecvSize() {
...
@@ -1620,24 +1612,3 @@ func (p *GrpcClientConfig) initClientMaxRecvSize() {
log
.
Debug
(
"initClientMaxRecvSize"
,
log
.
Debug
(
"initClientMaxRecvSize"
,
zap
.
String
(
"role"
,
p
.
Domain
),
zap
.
Int
(
"grpc.clientMaxRecvSize"
,
p
.
ClientMaxRecvSize
))
zap
.
String
(
"role"
,
p
.
Domain
),
zap
.
Int
(
"grpc.clientMaxRecvSize"
,
p
.
ClientMaxRecvSize
))
}
}
// CheckPortAvailable check if a port is available to be listened on
func
CheckPortAvailable
(
port
int
)
bool
{
addr
:=
":"
+
strconv
.
Itoa
(
port
)
listener
,
err
:=
net
.
Listen
(
"tcp"
,
addr
)
if
listener
!=
nil
{
listener
.
Close
()
}
return
err
==
nil
}
// GetAvailablePort return an available port that can be listened on
func
GetAvailablePort
()
int
{
listener
,
err
:=
net
.
Listen
(
"tcp"
,
":0"
)
if
err
!=
nil
{
panic
(
err
)
}
defer
listener
.
Close
()
return
listener
.
Addr
()
.
(
*
net
.
TCPAddr
)
.
Port
}
internal/util/paramtable/global_param_test.go
浏览文件 @
1601a61b
...
@@ -402,11 +402,3 @@ func TestGrpcClientParams(t *testing.T) {
...
@@ -402,11 +402,3 @@ func TestGrpcClientParams(t *testing.T) {
Params
.
initClientMaxSendSize
()
Params
.
initClientMaxSendSize
()
assert
.
Equal
(
t
,
Params
.
ClientMaxSendSize
,
DefaultClientMaxSendSize
)
assert
.
Equal
(
t
,
Params
.
ClientMaxSendSize
,
DefaultClientMaxSendSize
)
}
}
func
TestCheckPortAvailable
(
t
*
testing
.
T
)
{
num
:=
10
for
i
:=
0
;
i
<
num
;
i
++
{
port
:=
GetAvailablePort
()
assert
.
Equal
(
t
,
CheckPortAvailable
(
port
),
true
)
}
}
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录