Skip to content
Projects
Groups
Snippets
Help
Loading...
Help
Support
Submit feedback
Contribute to GitLab
Sign in
Toggle navigation
M
mogo
Project
Project
Details
Activity
Releases
Cycle Analytics
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Charts
Issues
0
Issues
0
List
Boards
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Charts
Wiki
Wiki
Snippets
Snippets
Members
Members
Collapse sidebar
Close sidebar
Activity
Graph
Charts
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
Odysseus
mogo
Commits
0b9e36dd
Commit
0b9e36dd
authored
Jun 28, 2024
by
vicotor
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
update workerinfo
parent
08aca3ed
Changes
2
Show whitespace changes
Inline
Side-by-side
Showing
2 changed files
with
159 additions
and
25 deletions
+159
-25
workerinfo.go
operator/workerinfo.go
+15
-2
workerinfo_test.go
operator/workerinfo_test.go
+144
-23
No files found.
operator/workerinfo.go
View file @
0b9e36dd
...
...
@@ -29,6 +29,19 @@ func NewDBWorker(client *mongo.Client, database string) *WorkerInfoOperator {
}
}
func
(
d
*
WorkerInfoOperator
)
CreateIndex
(
ctx
context
.
Context
)
error
{
_
,
err
:=
d
.
col
.
Indexes
()
.
CreateMany
(
ctx
,
[]
mongo
.
IndexModel
{
{
Keys
:
bson
.
D
{{
"worker_id"
,
1
}},
},
{
Keys
:
bson
.
D
{{
"model_infos.running_models.wait_time"
,
1
},
{
"model_infos.running_models.model_id"
,
1
}},
},
})
return
err
}
func
(
d
*
WorkerInfoOperator
)
Clear
()
{
d
.
col
.
DeleteMany
(
context
.
Background
(),
bson
.
M
{})
}
...
...
@@ -218,9 +231,9 @@ func (d *WorkerInfoOperator) FindWorkerByInstallModelAndSortByGpuRam(ctx context
// sort by gpu ram
findOptions
:=
options
.
Find
()
findOptions
.
SetLimit
(
int64
(
limit
))
findOptions
.
SetSort
(
bson
.
D
{{
"hardware.GPU.
ram
"
,
-
1
}})
findOptions
.
SetSort
(
bson
.
D
{{
"hardware.GPU.
mem_free
"
,
-
1
}})
selector
:=
bson
.
M
{
"model_infos.installed_models.model_id"
:
modelId
,
"hardware.GPU.performance"
:
bson
.
M
{
"$gte"
:
performance
},
"hardware.GPU.
ram
"
:
bson
.
M
{
"$gte"
:
ram
}}
selector
:=
bson
.
M
{
"model_infos.installed_models.model_id"
:
modelId
,
"hardware.GPU.performance"
:
bson
.
M
{
"$gte"
:
performance
},
"hardware.GPU.
mem_free
"
:
bson
.
M
{
"$gte"
:
ram
}}
cursor
,
err
:=
d
.
col
.
Find
(
ctx
,
selector
,
findOptions
)
if
err
!=
nil
{
return
nil
,
err
...
...
operator/workerinfo_test.go
View file @
0b9e36dd
...
...
@@ -13,14 +13,16 @@ import (
"log"
"math/rand"
"strconv"
"sync"
"testing"
"time"
)
var
(
maxModelId
=
1000
0
maxModelId
=
1000
idlist
=
make
([]
string
,
0
,
1000000
)
database
=
"test"
once
=
sync
.
Once
{}
)
func
ConnectMongoDB
()
(
*
mongo
.
Client
,
error
)
{
...
...
@@ -39,40 +41,51 @@ func ConnectMongoDB() (*mongo.Client, error) {
return
client
,
nil
}
//func init() {
// client, err := ConnectMongoDB()
// if err != nil {
// log.Fatal(err)
// }
// idlist = initdata(client)
//}
func
initdata
(
client
*
mongo
.
Client
)
[]
string
{
func
initdata
(
client
*
mongo
.
Client
,
count
int
,
running
bool
,
installed
bool
)
[]
string
{
t1
:=
time
.
Now
()
db
:=
NewDBWorker
(
client
,
database
)
dbRunning
:=
NewDBWorkerRunning
(
client
,
database
)
dbInstalled
:=
NewDBWorkerInstalled
(
client
,
database
)
// Insert 1,000,000 DbWorkerInfo to operator
for
i
:=
0
;
i
<
1000
;
i
++
{
for
i
:=
0
;
i
<
count
;
i
++
{
worker
:=
generateAWorker
()
result
,
err
:=
db
.
InsertWorker
(
context
.
Background
(),
worker
)
if
err
!=
nil
{
panic
(
fmt
.
Sprintf
(
"insert worker failed with err:%s"
,
err
))
}
{
if
running
{
// add worker running info to dbRunning
runnings
:=
make
([]
*
WorkerRunningInfo
,
0
,
len
(
worker
.
Models
.
RunningModels
))
for
_
,
model
:=
range
worker
.
Models
.
RunningModels
{
id
,
_
:=
strconv
.
Atoi
(
model
.
ModelID
)
running
Info
:=
&
WorkerRunningInfo
{
running
s
=
append
(
runnings
,
&
WorkerRunningInfo
{
WorkerId
:
worker
.
WorkerId
,
ModelId
:
id
,
ExecTime
:
model
.
ExecTime
,
})
}
_
,
err
:=
dbRunning
.
Insert
(
context
.
Background
(),
runningInfo
)
_
,
err
=
dbRunning
.
InsertMany
(
context
.
Background
(),
runnings
)
if
err
!=
nil
{
panic
(
fmt
.
Sprintf
(
"insert worker
failed with err:%s"
,
err
))
panic
(
fmt
.
Sprintf
(
"insert worker running
failed with err:%s"
,
err
))
}
}
if
installed
{
// add worker installed info to dbInstalled
installeds
:=
make
([]
*
WorkerInstalledInfo
,
0
,
len
(
worker
.
Models
.
InstalledModels
))
for
_
,
model
:=
range
worker
.
Models
.
InstalledModels
{
id
,
_
:=
strconv
.
Atoi
(
model
.
ModelID
)
installeds
=
append
(
installeds
,
&
WorkerInstalledInfo
{
WorkerId
:
worker
.
WorkerId
,
ModelId
:
id
,
GpuFree
:
1024
*
1024
*
1024
,
GpuSeq
:
0
,
})
}
_
,
err
=
dbInstalled
.
InsertMany
(
context
.
Background
(),
installeds
)
if
err
!=
nil
{
panic
(
fmt
.
Sprintf
(
"insert worker installed failed with err:%s"
,
err
))
}
}
id
,
ok
:=
result
.
InsertedID
.
(
primitive
.
ObjectID
)
...
...
@@ -223,6 +236,15 @@ func generateAInstallModel() *types.InstalledModel {
}
}
func
generateAInstallModelWithId
(
id
int
)
*
types
.
InstalledModel
{
return
&
types
.
InstalledModel
{
ModelID
:
strconv
.
Itoa
(
id
),
DiskSize
:
101
,
InstalledTime
:
time
.
Now
()
.
Unix
(),
LastRunTime
:
time
.
Now
()
.
Unix
(),
}
}
func
generateARunningModel
()
*
types
.
RunningModel
{
return
&
types
.
RunningModel
{
ModelID
:
getRandId
(
maxModelId
),
...
...
@@ -234,13 +256,24 @@ func generateARunningModel() *types.RunningModel {
ExecTime
:
rand
.
Intn
(
100
),
}
}
func
generateARunningModelWithId
(
id
int
)
*
types
.
RunningModel
{
return
&
types
.
RunningModel
{
ModelID
:
strconv
.
Itoa
(
id
),
GpuSeq
:
rand
.
Intn
(
3
),
GpuRAM
:
generateAGpuRam
(),
StartedTime
:
time
.
Now
()
.
Unix
(),
LastWorkTime
:
time
.
Now
()
.
Unix
(),
TotalRunCount
:
rand
.
Intn
(
100
),
ExecTime
:
rand
.
Intn
(
100
),
}
}
func
generateAModel
()
*
types
.
ModelInfo
{
m
:=
&
types
.
ModelInfo
{
InstalledModels
:
make
([]
*
types
.
InstalledModel
,
0
,
1000
),
RunningModels
:
make
([]
*
types
.
RunningModel
,
0
,
1000
),
}
for
i
:=
0
;
i
<
1
00
;
i
++
{
for
i
:=
0
;
i
<
2
00
;
i
++
{
m
.
InstalledModels
=
append
(
m
.
InstalledModels
,
generateAInstallModel
())
if
len
(
m
.
RunningModels
)
<
100
{
m
.
RunningModels
=
append
(
m
.
RunningModels
,
generateARunningModel
())
...
...
@@ -479,6 +512,13 @@ func BenchmarkDbWorker_InsertWorker(b *testing.B) {
log
.
Fatal
(
err
)
}
db
:=
NewDBWorker
(
client
,
database
)
once
.
Do
(
func
()
{
db
.
Clear
()
if
err
:=
db
.
CreateIndex
(
context
.
Background
());
err
!=
nil
{
panic
(
fmt
.
Sprintf
(
"create index failed with err:%s"
,
err
))
}
})
defer
db
.
client
.
Disconnect
(
context
.
Background
())
b
.
ResetTimer
()
for
i
:=
0
;
i
<
b
.
N
;
i
++
{
...
...
@@ -497,6 +537,12 @@ func BenchmarkDbWorker_InsertWorker_Parallel(b *testing.B) {
log
.
Fatal
(
err
)
}
db
:=
NewDBWorker
(
client
,
database
)
once
.
Do
(
func
()
{
db
.
Clear
()
if
err
:=
db
.
CreateIndex
(
context
.
Background
());
err
!=
nil
{
panic
(
fmt
.
Sprintf
(
"create index failed with err:%s"
,
err
))
}
})
defer
db
.
client
.
Disconnect
(
context
.
Background
())
b
.
RunParallel
(
func
(
pb
*
testing
.
PB
)
{
for
pb
.
Next
()
{
...
...
@@ -513,8 +559,18 @@ func BenchmarkDbWorker_UpdateHardware(b *testing.B) {
if
err
!=
nil
{
log
.
Fatal
(
err
)
}
b
.
StopTimer
()
db
:=
NewDBWorker
(
client
,
database
)
defer
db
.
client
.
Disconnect
(
context
.
Background
())
once
.
Do
(
func
()
{
db
.
Clear
()
if
err
:=
db
.
CreateIndex
(
context
.
Background
());
err
!=
nil
{
panic
(
fmt
.
Sprintf
(
"create index failed with err:%s"
,
err
))
}
idlist
=
initdata
(
client
,
10000
,
false
,
false
)
})
b
.
StartTimer
()
b
.
ResetTimer
()
for
i
:=
0
;
i
<
b
.
N
;
i
++
{
...
...
@@ -535,13 +591,26 @@ func BenchmarkDbWorker_UpdateHardware_Parallel(b *testing.B) {
if
err
!=
nil
{
log
.
Fatal
(
err
)
}
b
.
StopTimer
()
db
:=
NewDBWorker
(
client
,
database
)
defer
db
.
client
.
Disconnect
(
context
.
Background
())
once
.
Do
(
func
()
{
db
.
Clear
()
if
err
:=
db
.
CreateIndex
(
context
.
Background
());
err
!=
nil
{
panic
(
fmt
.
Sprintf
(
"create index failed with err:%s"
,
err
))
}
idlist
=
initdata
(
client
,
10000
,
false
,
false
)
})
b
.
StartTimer
()
b
.
ResetTimer
()
b
.
RunParallel
(
func
(
pb
*
testing
.
PB
)
{
for
pb
.
Next
()
{
b
.
StopTimer
()
idx
:=
rand
.
Intn
(
len
(
idlist
))
nhardware
:=
generateAHardware
()
b
.
StartTimer
()
if
err
:=
db
.
UpdateHardware
(
context
.
Background
(),
idlist
[
idx
],
nhardware
);
err
!=
nil
{
panic
(
fmt
.
Sprintf
(
"update worker failed with err:%s"
,
err
))
}
...
...
@@ -554,17 +623,29 @@ func BenchmarkDbWorker_FindWorkerByInstallModelAndSortByGpuRam(b *testing.B) {
if
err
!=
nil
{
log
.
Fatal
(
err
)
}
b
.
StopTimer
()
db
:=
NewDBWorker
(
client
,
database
)
defer
db
.
client
.
Disconnect
(
context
.
Background
())
once
.
Do
(
func
()
{
db
.
Clear
()
if
err
:=
db
.
CreateIndex
(
context
.
Background
());
err
!=
nil
{
panic
(
fmt
.
Sprintf
(
"create index failed with err:%s"
,
err
))
}
idlist
=
initdata
(
client
,
10000
,
false
,
false
)
})
b
.
StartTimer
()
b
.
ResetTimer
()
for
i
:=
0
;
i
<
b
.
N
;
i
++
{
b
.
StopTimer
()
installedModelId
:=
getRandId
(
maxModelId
)
performance
:=
generateAGpuPerformance
()
ram
:=
generateAGpuRam
()
ram
:=
int64
(
100
*
1024
)
//generateAGpuRam()
b
.
StartTimer
()
if
w
,
err
:=
db
.
FindWorkerByInstallModelAndSortByGpuRam
(
context
.
Background
(),
installedModelId
,
performance
,
ram
,
10
);
err
!=
nil
{
panic
(
fmt
.
Sprintf
(
"find worker failed with err:%s"
,
err
))
}
else
if
len
(
w
)
==
0
{
b
.
Logf
(
"FindWorkerByInstallModelAndSortByGpuRam find %d with id %s
\n
"
,
len
(
w
),
installedModelId
)
//
b.Logf("FindWorkerByInstallModelAndSortByGpuRam find %d with id %s\n", len(w), installedModelId)
}
}
}
...
...
@@ -574,17 +655,32 @@ func BenchmarkDbWorker_FindWorkerByInstallModelAndSortByGpuRam_Parallel(b *testi
if
err
!=
nil
{
log
.
Fatal
(
err
)
}
b
.
StopTimer
()
db
:=
NewDBWorker
(
client
,
database
)
defer
db
.
client
.
Disconnect
(
context
.
Background
())
once
.
Do
(
func
()
{
db
.
Clear
()
if
err
:=
db
.
CreateIndex
(
context
.
Background
());
err
!=
nil
{
panic
(
fmt
.
Sprintf
(
"create index failed with err:%s"
,
err
))
}
idlist
=
initdata
(
client
,
10000
,
false
,
false
)
})
b
.
StartTimer
()
b
.
ResetTimer
()
b
.
RunParallel
(
func
(
pb
*
testing
.
PB
)
{
for
pb
.
Next
()
{
b
.
StopTimer
()
installedModelId
:=
getRandId
(
maxModelId
)
performance
:=
generateAGpuPerformance
()
ram
:=
generateAGpuRam
()
ram
:=
int64
(
100
*
1024
)
//generateAGpuRam()
b
.
StartTimer
()
if
w
,
err
:=
db
.
FindWorkerByInstallModelAndSortByGpuRam
(
context
.
Background
(),
installedModelId
,
performance
,
ram
,
10
);
err
!=
nil
{
panic
(
fmt
.
Sprintf
(
"find worker failed with err:%s"
,
err
))
}
else
if
len
(
w
)
==
0
{
b
.
Logf
(
"FindWorkerByInstallModelAndSortByGpuRam find %d with id %s
\n
"
,
len
(
w
),
installedModelId
)
//b.Logf("FindWorkerByInstallModelAndSortByGpuRam find %d with id %s\n", len(w), installedModelId)
}
else
{
}
}
})
...
...
@@ -595,11 +691,23 @@ func BenchmarkDbWorker_FindWorkerByRunningModelAndSortByWaitTime(b *testing.B) {
if
err
!=
nil
{
log
.
Fatal
(
err
)
}
b
.
StopTimer
()
db
:=
NewDBWorker
(
client
,
database
)
defer
db
.
client
.
Disconnect
(
context
.
Background
())
once
.
Do
(
func
()
{
db
.
Clear
()
if
err
:=
db
.
CreateIndex
(
context
.
Background
());
err
!=
nil
{
panic
(
fmt
.
Sprintf
(
"create index failed with err:%s"
,
err
))
}
idlist
=
initdata
(
client
,
1000
,
false
,
false
)
})
b
.
StartTimer
()
b
.
ResetTimer
()
for
i
:=
0
;
i
<
b
.
N
;
i
++
{
b
.
StopTimer
()
runningModelId
:=
getRandId
(
maxModelId
)
b
.
StartTimer
()
if
w
,
err
:=
db
.
FindWorkerByRunningModelAndSortByWaitTime
(
context
.
Background
(),
runningModelId
,
10
);
err
!=
nil
{
panic
(
fmt
.
Sprintf
(
"find worker failed with err:%s"
,
err
))
}
else
if
len
(
w
)
==
0
{
...
...
@@ -613,11 +721,24 @@ func BenchmarkDbWorker_FindWorkerByRunningModelAndSortByWaitTime_Parallel(b *tes
if
err
!=
nil
{
log
.
Fatal
(
err
)
}
b
.
StopTimer
()
db
:=
NewDBWorker
(
client
,
database
)
defer
db
.
client
.
Disconnect
(
context
.
Background
())
once
.
Do
(
func
()
{
db
.
Clear
()
if
err
:=
db
.
CreateIndex
(
context
.
Background
());
err
!=
nil
{
panic
(
fmt
.
Sprintf
(
"create index failed with err:%s"
,
err
))
}
idlist
=
initdata
(
client
,
1000
,
false
,
false
)
})
b
.
StartTimer
()
b
.
ResetTimer
()
b
.
RunParallel
(
func
(
pb
*
testing
.
PB
)
{
for
pb
.
Next
()
{
b
.
StopTimer
()
runningModelId
:=
getRandId
(
maxModelId
)
b
.
StartTimer
()
if
w
,
err
:=
db
.
FindWorkerByRunningModelAndSortByWaitTime
(
context
.
Background
(),
runningModelId
,
10
);
err
!=
nil
{
panic
(
fmt
.
Sprintf
(
"find worker failed with err:%s"
,
err
))
}
else
if
len
(
w
)
==
0
{
...
...
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment