@@ -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
0 commit comments