Skip to content

Commit 43a69da

Browse files
Chen Zhilingfeast-ci-bot
authored andcommitted
Fix batch client call missing request, fix multi-entity lookup (#299)
1 parent b90facd commit 43a69da

2 files changed

Lines changed: 46 additions & 22 deletions

File tree

sdk/python/feast/client.py

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -349,7 +349,8 @@ def get_batch_features(
349349

350350
# Retrieve serving information to determine store type and staging location
351351
serving_info = (
352-
self._serving_service_stub.GetFeastServingInfo()
352+
self._serving_service_stub.GetFeastServingInfo(
353+
GetFeastServingInfoRequest(), timeout=GRPC_CONNECTION_TIMEOUT_DEFAULT)
353354
) # type: GetFeastServingInfoResponse
354355

355356
if serving_info.type != FeastServingType.FEAST_SERVING_TYPE_BATCH:
Lines changed: 44 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,4 @@
1-
SELECT {{ fullEntitiesList | join(', ')}}, {% for featureSet in featureSets %}{% for featureName in featureSet.features %}{{ featureSet.name }}_v{{ featureSet.version }}_{{ featureName }}, {% endfor %}{% endfor %}event_timestamp FROM (WITH union_features AS (
2-
SELECT
1+
WITH union_features AS (SELECT
32
event_timestamp,
43
{% for featureSet in featureSets %}NULL as {{ featureSet.name }}_v{{ featureSet.version }}_feature_timestamp,
54
{% endfor %}{{ fullEntitiesList | join(', ')}},
@@ -24,38 +23,62 @@ SELECT
2423
{% endif %}
2524
{% endfor %}
2625
false AS is_entity_table
27-
FROM `{{projectId}}.{{datasetId}}.{{ featureSet.name }}_v{{ featureSet.version }}` WHERE event_timestamp <= '{{maxTimestamp}}' AND event_timestamp >= '{{minTimestamp}}'
26+
FROM `{{projectId}}.{{datasetId}}.{{ featureSet.name }}_v{{ featureSet.version }}` WHERE event_timestamp <= '{{maxTimestamp}}' AND event_timestamp >= Timestamp_sub(TIMESTAMP '{{ minTimestamp }}', interval {{ featureSet.maxAge }} second)
2827
{% endfor %}
29-
)
30-
SELECT
31-
event_timestamp,
32-
{% for featureSet in featureSets %}
33-
IF(event_timestamp >= {{ featureSet.name }}_v{{ featureSet.version }}_feature_timestamp AND Timestamp_sub(event_timestamp, interval {{ featureSet.maxAge }} second) < {{ featureSet.name }}_v{{ featureSet.version }}_feature_timestamp, {{ featureSet.name }}_v{{ featureSet.version }}_feature_timestamp, NULL) as {{ featureSet.name }}_v{{ featureSet.version }}_feature_timestamp,
34-
{% endfor %}
35-
{{ fullEntitiesList | join(', ')}}
36-
FROM (
28+
), ts_union AS (
29+
{% for featureSet in featureSets %}
30+
SELECT * FROM (
31+
SELECT
32+
event_timestamp,
33+
{{ fullEntitiesList | join(', ')}},
34+
{% for otherFeatureSet in featureSets %}
35+
{% if otherFeatureSet.id == featureSet.id %}
36+
LAST_VALUE({{ otherFeatureSet.name }}_v{{ otherFeatureSet.version }}_feature_timestamp IGNORE NULLS) over w AS {{ otherFeatureSet.name }}_v{{ otherFeatureSet.version }}_feature_timestamp,
37+
{% else %}
38+
{{ otherFeatureSet.name }}_v{{ otherFeatureSet.version }}_feature_timestamp,
39+
{% endif %}
40+
{% endfor %}
41+
is_entity_table
42+
FROM union_features
43+
WINDOW w AS (PARTITION BY {{ featureSet.entities | join(', ') }} ORDER BY event_timestamp, is_entity_table ASC ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW)
44+
) WHERE is_entity_table
45+
{% if loop.last %}
46+
{% else %}
47+
UNION ALL
48+
{% endif %}
49+
{% endfor %}
50+
), ts_coalesce AS (
3751
SELECT
3852
event_timestamp,
3953
{{ fullEntitiesList | join(', ')}},
4054
{% for featureSet in featureSets %}
4155
LAST_VALUE({{ featureSet.name }}_v{{ featureSet.version }}_feature_timestamp IGNORE NULLS) over w AS {{ featureSet.name }}_v{{ featureSet.version }}_feature_timestamp,
4256
{% endfor %}
43-
is_entity_table
44-
FROM union_features
45-
WINDOW w AS (PARTITION BY source ORDER BY event_timestamp ASC)
46-
) WHERE is_entity_table)
57+
ROW_NUMBER() over w as rn
58+
FROM ts_union
59+
WINDOW w AS (PARTITION BY {{ fullEntitiesList | join(', ')}}, event_timestamp ORDER BY event_timestamp)
60+
), ts_final AS (
61+
SELECT
62+
event_timestamp,
63+
{{ fullEntitiesList | join(', ')}},
64+
{% for featureSet in featureSets %}
65+
IF(event_timestamp >= {{ featureSet.name }}_v{{ featureSet.version }}_feature_timestamp AND Timestamp_sub(event_timestamp, interval {{ featureSet.maxAge }} second) < {{ featureSet.name }}_v{{ featureSet.version }}_feature_timestamp, {{ featureSet.name }}_v{{ featureSet.version }}_feature_timestamp, NULL) as {{ featureSet.name }}_v{{ featureSet.version }}_feature_timestamp{% if loop.last %}{% else %}, {% endif %}
66+
{% endfor %}
67+
FROM ts_coalesce WHERE rn = 1
68+
)
69+
SELECT * FROM ts_final
4770
{% for featureSet in featureSets %}
48-
LEFT JOIN
49-
(SELECT {{ featureSet.name }}_v{{ featureSet.version }}_feature_timestamp,
71+
LEFT JOIN
72+
(SELECT {{ featureSet.name }}_v{{ featureSet.version }}_feature_timestamp,
5073
{% for featureName in featureSet.features %}{{ featureSet.name }}_v{{ featureSet.version }}_{{ featureName }},
5174
{% endfor %}{{ featureSet.entities | join(', ') }}
52-
FROM (SELECT
75+
FROM (SELECT
5376
event_timestamp as {{ featureSet.name }}_v{{ featureSet.version }}_feature_timestamp,
5477
{{ featureSet.entities | join(', ') }},
5578
{% for featureName in featureSet.features %}
5679
{{ featureName }} as {{ featureSet.name }}_v{{ featureSet.version }}_{{ featureName }},
5780
{% endfor %}ROW_NUMBER() OVER(PARTITION BY event_timestamp, {{ featureSet.entities | join(', ') }} ORDER BY created_timestamp DESC) as {{ featureSet.name }}_v{{ featureSet.version }}_rown
58-
FROM `{{ projectId }}.{{ datasetId }}.{{ featureSet.name }}_v{{ featureSet.version }}` WHERE event_timestamp <= '{{maxTimestamp}}' AND event_timestamp >= '{{minTimestamp}}'
81+
FROM `{{ projectId }}.{{ datasetId }}.{{ featureSet.name }}_v{{ featureSet.version }}` WHERE event_timestamp <= '{{maxTimestamp}}' AND event_timestamp >= Timestamp_sub(TIMESTAMP '{{ minTimestamp }}', interval {{ featureSet.maxAge }} second)
5982
) WHERE {{ featureSet.name }}_v{{ featureSet.version }}_rown = 1
60-
) USING ({{ featureSet.name }}_v{{ featureSet.version }}_feature_timestamp, source)
61-
{% endfor %}
83+
) USING ({{ featureSet.name }}_v{{ featureSet.version }}_feature_timestamp, {{ featureSet.entities | join(', ') }})
84+
{% endfor %}

0 commit comments

Comments
 (0)