Skip to content

Commit a77a73b

Browse files
Ly Caoachals
authored andcommitted
added goroutines to OnlineRead
Signed-off-by: Felix Wang <wangfelix98@gmail.com> Signed-off-by: Achal Shah <achals@gmail.com>
1 parent 6887787 commit a77a73b

4 files changed

Lines changed: 89 additions & 45 deletions

File tree

go/feast/featurestore.go

Lines changed: 77 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -220,20 +220,87 @@ func (fs *FeatureStore) GetOnlineFeatures(request *serving.GetOnlineFeaturesRequ
220220
joinKeyValues[DUMMY_ENTITY_ID] = dummyEntityColumn
221221
}
222222

223-
for table, requestedFeatures := range groupedRefs {
224-
tableEntityValues, idxs, err := fs.getUniqueEntities( table, joinKeyValues, entityNameToJoinKeyMap)
225-
if err != nil {
226-
return nil, err
223+
/*
224+
TODO (Ly): Precompute indices that need to be written by each table
225+
and have separate goroutines for each write. For simplicity, I obmit
226+
this optimization for now
227+
*/
228+
229+
230+
// Writer
231+
type featureDataWriter struct {
232+
featureData [][]FeatureData
233+
idxs [][]int
234+
requestedFeatures []string
235+
table *FeatureView
236+
}
237+
// Write one new read at a time
238+
featureDataWriterChan := make(chan *featureDataWriter, 1)
239+
errorFromReadChan := make(chan error, 1)
240+
writeCount := len(groupedRefs)
241+
stopReadChans := make([]chan struct{}, writeCount)
242+
writeDone := make(chan struct{}, 1)
243+
for index = 0; index < writeCount; index++ {
244+
stopReadChans[index] = make(chan struct{}, 1)
245+
}
246+
go func() {
247+
wrote := 0
248+
var featureDataWrite *featureDataWriter
249+
for wrote < writeCount {
250+
featureDataWrite = <- featureDataWriterChan
251+
fs.populateResponseFromFeatureData( featureDataWrite.featureData,
252+
featureDataWrite.idxs,
253+
onlineFeatureResponse,
254+
fullFeatureNames,
255+
featureDataWrite.requestedFeatures,
256+
featureDataWrite.table,
257+
)
258+
wrote += 1
227259
}
228-
featureData, err := fs.readFromOnlineStore(tableEntityValues, requestedFeatures, table)
229-
if err != nil {
260+
// TODO (Ly): ODFV, skip augmentResponseWithOnDemandTransforms
261+
fs.dropUnneededColumns(onlineFeatureResponse, requestedResultRowNames)
262+
writeDone <- struct{}{}
263+
}()
264+
265+
index = 0
266+
for table, requestedFeatures := range groupedRefs {
267+
select {
268+
case err := <- errorFromReadChan:
269+
for index = 0; index < writeCount; index++ {
270+
stopReadChans[index] <- struct{}{}
271+
}
230272
return nil, err
273+
default:
274+
275+
stopReadChans[index] = make(chan struct{}, 1)
276+
go func(stopChan chan struct{}, table *FeatureView, requestedFeatures []string) {
277+
select {
278+
case <- stopChan:
279+
default:
280+
tableEntityValues, idxs, err := fs.getUniqueEntities( table, joinKeyValues, entityNameToJoinKeyMap)
281+
if err != nil {
282+
select {
283+
case errorFromReadChan <- err:
284+
default:
285+
}
286+
return
287+
}
288+
featureData, err := fs.readFromOnlineStore(tableEntityValues, requestedFeatures, table)
289+
if err != nil {
290+
// Don't block if another read failed in the middle of read
291+
select {
292+
case errorFromReadChan <- err:
293+
default:
294+
}
295+
} else {
296+
featureDataWriterChan <- &featureDataWriter{ featureData, idxs, requestedFeatures, table }
297+
}
298+
}
299+
}(stopReadChans[index], table, requestedFeatures)
300+
index += 1
231301
}
232-
fs.populateResponseFromFeatureData(featureData, idxs, onlineFeatureResponse, fullFeatureNames, requestedFeatures, table)
233302
}
234-
235-
// TODO (Ly): ODFV, skip augmentResponseWithOnDemandTransforms
236-
fs.dropUnneededColumns(onlineFeatureResponse, requestedResultRowNames)
303+
<- writeDone
237304
return onlineFeatureResponse, nil
238305
}
239306

go/server/main.go

Lines changed: 0 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -36,7 +36,6 @@ func main() {
3636
log.Fatalln(fmt.Sprintf("One of %s of %s environment variables must be set", flagFeastRepoPath, flagFeastRepoConfig))
3737
return
3838
}
39-
log.Println(repoPath, repoConfig)
4039
config, err := feast.NewRepoConfig(repoPath, repoConfig)
4140
if err != nil {
4241
log.Fatalln(err)
@@ -55,43 +54,24 @@ func main() {
5554
grpcPort = defaultFeastGrpcPort
5655
}
5756

58-
// fmt.Println("starting for loop")
5957
reader := bufio.NewReader(os.Stdin)
6058
writer := bufio.NewWriter(os.Stdout)
6159
serverCounts := 0
6260
for {
6361

6462
text, _ := reader.ReadString('\n')
6563
text = strings.Trim(text, "\n")
66-
// fmt.Println("Received from stdin", text)
6764
commands := strings.Split(text, " ")
6865
if len(commands) == 0 {
6966
log.Fatalln(errors.New("Invalid command. Should be [startGrpc] or [startHttp host:port] or [stop]"))
7067
return
7168
} else if commands[0] == "startGrpc" {
72-
// fmt.Fprintf(writer, "Success!")
7369
writer.Flush()
7470
go func() {
7571
startGrpcServer(fs, grpcPort)
76-
// server := servingServiceServer{
77-
// fs: fs,
78-
// }
79-
// log.Printf("Starting a gRPC server at port %s...", grpcPort)
80-
// lis, err := net.Listen("tcp", fmt.Sprintf(":%s", grpcPort))
81-
// if err != nil {
82-
// log.Fatalln(err)
83-
// }
84-
// grpcServer := grpc.NewServer()
85-
// defer grpcServer.Stop()
86-
// serving.RegisterServingServiceServer(grpcServer, &server)
87-
// err = grpcServer.Serve(lis)
88-
// if err != nil {
89-
// log.Fatalln(err)
90-
// }
9172
}()
9273
serverCounts += 1
9374
} else if commands[0] == "startHttp" {
94-
// fmt.Fprintf(writer, "Success!")
9575
writer.Flush()
9676
if len(commands) < 2 {
9777
log.Fatalln(errors.New("Invalid command. Should be: startHttp host:port"))
@@ -109,7 +89,6 @@ func main() {
10989
}
11090

11191
func startGrpcServer(fs *feast.FeatureStore, grpcPort string) {
112-
fmt.Println("startGrpcServer", grpcPort)
11392
server := servingServiceServer{
11493
fs: fs,
11594
}
@@ -130,7 +109,6 @@ func startGrpcServer(fs *feast.FeatureStore, grpcPort string) {
130109
}
131110

132111
func startHttpServer(fs *feast.FeatureStore, address string) {
133-
fmt.Println("startHttpServer", address)
134112
http.HandleFunc("/get-online-features", func(w http.ResponseWriter, req *http.Request) {
135113
reqBodyBytes, err := io.ReadAll(req.Body)
136114
var grpcRequest serving.GetOnlineFeaturesRequest

sdk/python/feast/go_server.py

Lines changed: 11 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -139,7 +139,7 @@ def start_grpc_server(self):
139139
self.stop()
140140
raise errors.GoSubprocessConnectionFailed() from e
141141
# Sleep for 0.1 second before retrying
142-
time.sleep(0.1)
142+
time.sleep(0.1)
143143

144144
def start_http_server(self, host: str, port: int):
145145
if self.httpServerStarted:
@@ -162,16 +162,16 @@ def stop(self):
162162
# Only send sigkill if there's a problem telling go subprocess to stop
163163
# Otherwise, let go subprocess clean up and shut down itself
164164
if not self.pipeClosed:
165-
for i in range(10):
166-
try:
167-
self.process.stdin.write(bytes(f"stop\n", encoding='utf8'))
168-
self.process.stdin.flush()
169-
# time.sleep(0.1)
170-
# self.process.stdin.close()
171-
break
172-
except subprocess.CalledProcessError as error:
173-
self.process.terminate()
174-
raise errors.GoSubprocessConnectionFailed() from error
165+
try:
166+
self.process.stdin.write(bytes(f"stop\n", encoding='utf8'))
167+
self.process.stdin.flush()
168+
# TODO (Ly): Review: We don't close stdin here
169+
# since if the call succeeds go process closes
170+
# itself and stdin?
171+
# self.process.stdin.close()
172+
except subprocess.CalledProcessError as error:
173+
self.process.terminate()
174+
raise errors.GoSubprocessConnectionFailed() from error
175175

176176
self.grpcServerStarted = False
177177
self.httpServerStarted = False

sdk/python/tests/conftest.py

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -199,8 +199,7 @@ def go_server_environment(request, worker_id: str):
199199
time.sleep(3)
200200

201201
def cleanup():
202-
if e.feature_store:
203-
e.feature_store.stop_go_server()
202+
e.feature_store.stop_go_server()
204203

205204
e.feature_store.teardown()
206205
if proc.is_alive():

0 commit comments

Comments
 (0)