Skip to content

Commit 6e08634

Browse files
committed
Fix tests
Signed-off-by: Terence <terencelimxp@gmail.com>
1 parent cd41e28 commit 6e08634

3 files changed

Lines changed: 28 additions & 30 deletions

File tree

tests/e2e/test_online_features.py

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -127,7 +127,7 @@ def test_streaming_ingestion(
127127

128128
if not pytestconfig.getoption("scheduled_streaming_job"):
129129
job = feast_client.start_stream_to_online_ingestion(feature_table)
130-
assert job.get_feature_table()[:63] == feature_table.name[:63]
130+
assert job.get_feature_table() == feature_table.name
131131
wait_retry_backoff(
132132
lambda: (None, job.get_status() == SparkJobStatus.IN_PROGRESS), 180
133133
)
@@ -172,7 +172,7 @@ def ingest_and_verify(
172172
original.event_timestamp.min().to_pydatetime(),
173173
original.event_timestamp.max().to_pydatetime() + timedelta(seconds=1),
174174
)
175-
assert job.get_feature_table()[:63] == feature_table.name[:63]
175+
assert job.get_feature_table() == feature_table.name
176176

177177
wait_retry_backoff(
178178
lambda: (None, job.get_status() == SparkJobStatus.COMPLETED), 180
@@ -200,7 +200,6 @@ def test_list_jobs_long_table_name(feast_client: Client, kafka_server, pytestcon
200200
kafka_broker,
201201
topic_name,
202202
"just1a2featuretable3with4a5really6really7really8really9really10really11really12long13name",
203-
avro_schema(),
204203
)
205204

206205
feast_client.apply(entity)

tests/e2e/test_validation.py

Lines changed: 3 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,3 @@
1-
import json
21
import time
32
import uuid
43

@@ -11,7 +10,7 @@
1110
from feast.contrib.validation.ge import apply_validation, create_validation_udf
1211
from feast.wait import wait_retry_backoff
1312
from tests.e2e.fixtures.statsd_stub import StatsDServer
14-
from tests.e2e.utils.common import create_schema, start_job, stop_job
13+
from tests.e2e.utils.common import avro_schema, create_schema, start_job, stop_job
1514
from tests.e2e.utils.kafka import check_consumer_exist, ingest_and_retrieve
1615

1716

@@ -39,9 +38,7 @@ def test_validation_with_ge(feast_client: Client, kafka_server, pytestconfig):
3938
kafka_broker = f"{kafka_server[0]}:{kafka_server[1]}"
4039
topic_name = f"avro-{uuid.uuid4()}"
4140

42-
entity, feature_table = create_schema(
43-
kafka_broker, topic_name, "validation_ge", avro_schema()
44-
)
41+
entity, feature_table = create_schema(kafka_broker, topic_name, "validation_ge")
4542
feast_client.apply_entity(entity)
4643
feast_client.apply_feature_table(feature_table)
4744

@@ -107,7 +104,7 @@ def test_validation_reports_metrics(
107104
topic_name = f"avro-{uuid.uuid4()}"
108105

109106
entity, feature_table = create_schema(
110-
kafka_broker, topic_name, "validation_ge_metrics", avro_schema()
107+
kafka_broker, topic_name, "validation_ge_metrics"
111108
)
112109
feast_client.apply_entity(entity)
113110
feast_client.apply_feature_table(feature_table)
@@ -181,21 +178,3 @@ def test_validation_reports_metrics(
181178
+ "\n"
182179
"Actual received metrics" + str(statsd_server.metrics),
183180
)
184-
185-
186-
def avro_schema():
187-
return json.dumps(
188-
{
189-
"type": "record",
190-
"name": "TestMessage",
191-
"fields": [
192-
{"name": "key", "type": "long"},
193-
{"name": "num", "type": "long"},
194-
{"name": "set", "type": "string"},
195-
{
196-
"name": "event_timestamp",
197-
"type": {"type": "long", "logicalType": "timestamp-micros"},
198-
},
199-
],
200-
}
201-
)

tests/e2e/utils/common.py

Lines changed: 23 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,5 @@
1+
import json
2+
13
from feast import (
24
Client,
35
Entity,
@@ -12,7 +14,7 @@
1214
from feast.wait import wait_retry_backoff
1315

1416

15-
def create_schema(kafka_broker, topic_name, feature_table_name, avro_schema):
17+
def create_schema(kafka_broker, topic_name, feature_table_name):
1618
entity = Entity(name="key", description="Key", value_type=ValueType.INT64)
1719
feature_table = FeatureTable(
1820
name=feature_table_name,
@@ -26,7 +28,7 @@ def create_schema(kafka_broker, topic_name, feature_table_name, avro_schema):
2628
stream_source=KafkaSource(
2729
event_timestamp_column="event_timestamp",
2830
bootstrap_servers=kafka_broker,
29-
message_format=AvroFormat(avro_schema),
31+
message_format=AvroFormat(avro_schema()),
3032
topic=topic_name,
3133
),
3234
)
@@ -39,7 +41,7 @@ def start_job(feast_client: Client, feature_table: FeatureTable, pytestconfig):
3941

4042
job = feast_client.start_stream_to_online_ingestion(feature_table)
4143
wait_retry_backoff(
42-
lambda: (None, job.get_status() == SparkJobStatus.IN_PROGRESS), 600
44+
lambda: (None, job.get_status() == SparkJobStatus.IN_PROGRESS), 180
4345
)
4446
return job
4547

@@ -49,3 +51,21 @@ def stop_job(job, feast_client: Client, feature_table: FeatureTable):
4951
job.cancel()
5052
else:
5153
feast_client.delete_feature_table(feature_table.name)
54+
55+
56+
def avro_schema():
57+
return json.dumps(
58+
{
59+
"type": "record",
60+
"name": "TestMessage",
61+
"fields": [
62+
{"name": "key", "type": "long"},
63+
{"name": "num", "type": "long"},
64+
{"name": "set", "type": "string"},
65+
{
66+
"name": "event_timestamp",
67+
"type": {"type": "long", "logicalType": "timestamp-micros"},
68+
},
69+
],
70+
}
71+
)

0 commit comments

Comments
 (0)