Skip to content

Commit b704dba

Browse files
committed
Refactor tests
Signed-off-by: Terence <terencelimxp@gmail.com>
1 parent 94c714a commit b704dba

4 files changed

Lines changed: 86 additions & 83 deletions

File tree

sdk/python/feast/pyspark/launchers/gcloud/dataproc.py

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -315,9 +315,9 @@ def dataproc_submit(
315315
}
316316

317317
if isinstance(job_params, StreamIngestionJobParameters):
318-
job_config["labels"][
319-
self.FEATURE_TABLE_LABEL_KEY
320-
] = job_params.get_feature_table_name()
318+
job_config["labels"][self.FEATURE_TABLE_LABEL_KEY] = _truncate_label(
319+
job_params.get_feature_table_name()
320+
)
321321
# Add job hash to labels only for the stream ingestion job
322322
job_config["labels"][self.JOB_HASH_LABEL_KEY] = job_params.get_job_hash()
323323

tests/e2e/test_online_features.py

Lines changed: 26 additions & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@
2323
from feast.data_format import AvroFormat, ParquetFormat
2424
from feast.pyspark.abc import SparkJobStatus
2525
from feast.wait import wait_retry_backoff
26+
from tests.e2e.utils.common import create_schema, start_job, stop_job
2627
from tests.e2e.utils.kafka import check_consumer_exist, ingest_and_retrieve
2728

2829

@@ -59,27 +60,6 @@ def test_offline_ingestion(
5960
ingest_and_verify(feast_client, feature_table, original)
6061

6162

62-
def test_offline_ingestion_long_table_name(
63-
feast_client: Client, batch_source: Union[BigQuerySource, FileSource]
64-
):
65-
entity = Entity(name="s2id", description="S2id", value_type=ValueType.INT64,)
66-
67-
feature_table = FeatureTable(
68-
name="just1a2featuretable3with4a5really6really7really8really9really10really11really12long13name",
69-
entities=["s2id"],
70-
features=[Feature("unique_drivers", ValueType.INT64)],
71-
batch_source=batch_source,
72-
)
73-
74-
feast_client.apply(entity)
75-
feast_client.apply(feature_table)
76-
77-
original = generate_data()
78-
feast_client.ingest(feature_table, original) # write to batch (offline) storage
79-
80-
ingest_and_verify(feast_client, feature_table, original)
81-
82-
8363
@pytest.mark.env("gcloud")
8464
def test_offline_ingestion_from_bq_view(pytestconfig, bq_dataset, feast_client: Client):
8565
original = generate_data()
@@ -193,13 +173,6 @@ def ingest_and_verify(
193173
original.event_timestamp.max().to_pydatetime() + timedelta(seconds=1),
194174
)
195175
assert job.get_feature_table()[:63] == feature_table.name[:63]
196-
all_job_ids = [
197-
job.get_id()
198-
for job in feast_client.list_jobs(
199-
include_terminated=True, table_name=feature_table.name
200-
)
201-
]
202-
assert job.get_id() in all_job_ids
203176

204177
wait_retry_backoff(
205178
lambda: (None, job.get_status() == SparkJobStatus.COMPLETED), 180
@@ -219,6 +192,31 @@ def ingest_and_verify(
219192
)
220193

221194

195+
def test_list_jobs_long_table_name(feast_client: Client, kafka_server, pytestconfig):
196+
kafka_broker = f"{kafka_server[0]}:{kafka_server[1]}"
197+
topic_name = f"avro-{uuid.uuid4()}"
198+
199+
entity, feature_table = create_schema(
200+
kafka_broker,
201+
topic_name,
202+
"just1a2featuretable3with4a5really6really7really8really9really10really11really12long13name",
203+
avro_schema(),
204+
)
205+
206+
feast_client.apply(entity)
207+
feast_client.apply(feature_table)
208+
209+
job = start_job(feast_client, feature_table, pytestconfig)
210+
all_job_ids = [
211+
job.get_id()
212+
for job in feast_client.list_jobs(
213+
include_terminated=True, table_name=feature_table.name
214+
)
215+
]
216+
assert job.get_id() in all_job_ids
217+
stop_job(job, feast_client, feature_table)
218+
219+
222220
def avro_schema():
223221
return json.dumps(
224222
{

tests/e2e/test_validation.py

Lines changed: 6 additions & 52 deletions
Original file line numberDiff line numberDiff line change
@@ -7,20 +7,11 @@
77
import pytest
88
from great_expectations.dataset import PandasDataset
99

10-
from feast import (
11-
Client,
12-
Entity,
13-
Feature,
14-
FeatureTable,
15-
FileSource,
16-
KafkaSource,
17-
ValueType,
18-
)
10+
from feast import Client
1911
from feast.contrib.validation.ge import apply_validation, create_validation_udf
20-
from feast.data_format import AvroFormat, ParquetFormat
21-
from feast.pyspark.abc import SparkJobStatus
2212
from feast.wait import wait_retry_backoff
2313
from tests.e2e.fixtures.statsd_stub import StatsDServer
14+
from tests.e2e.utils.common import create_schema, start_job, stop_job
2415
from tests.e2e.utils.kafka import check_consumer_exist, ingest_and_retrieve
2516

2617

@@ -44,50 +35,13 @@ def generate_test_data():
4435
return df
4536

4637

47-
def create_schema(kafka_broker, topic_name, feature_table_name):
48-
entity = Entity(name="key", description="Key", value_type=ValueType.INT64)
49-
feature_table = FeatureTable(
50-
name=feature_table_name,
51-
entities=["key"],
52-
features=[Feature("num", ValueType.INT64), Feature("set", ValueType.STRING)],
53-
batch_source=FileSource(
54-
event_timestamp_column="event_timestamp",
55-
file_format=ParquetFormat(),
56-
file_url="/dev/null",
57-
),
58-
stream_source=KafkaSource(
59-
event_timestamp_column="event_timestamp",
60-
bootstrap_servers=kafka_broker,
61-
message_format=AvroFormat(avro_schema()),
62-
topic=topic_name,
63-
),
64-
)
65-
return entity, feature_table
66-
67-
68-
def start_job(feast_client: Client, feature_table: FeatureTable, pytestconfig):
69-
if pytestconfig.getoption("scheduled_streaming_job"):
70-
return
71-
72-
job = feast_client.start_stream_to_online_ingestion(feature_table)
73-
wait_retry_backoff(
74-
lambda: (None, job.get_status() == SparkJobStatus.IN_PROGRESS), 120
75-
)
76-
return job
77-
78-
79-
def stop_job(job, feast_client: Client, feature_table: FeatureTable):
80-
if job:
81-
job.cancel()
82-
else:
83-
feast_client.delete_feature_table(feature_table.name)
84-
85-
8638
def test_validation_with_ge(feast_client: Client, kafka_server, pytestconfig):
8739
kafka_broker = f"{kafka_server[0]}:{kafka_server[1]}"
8840
topic_name = f"avro-{uuid.uuid4()}"
8941

90-
entity, feature_table = create_schema(kafka_broker, topic_name, "validation_ge")
42+
entity, feature_table = create_schema(
43+
kafka_broker, topic_name, "validation_ge", avro_schema()
44+
)
9145
feast_client.apply_entity(entity)
9246
feast_client.apply_feature_table(feature_table)
9347

@@ -153,7 +107,7 @@ def test_validation_reports_metrics(
153107
topic_name = f"avro-{uuid.uuid4()}"
154108

155109
entity, feature_table = create_schema(
156-
kafka_broker, topic_name, "validation_ge_metrics"
110+
kafka_broker, topic_name, "validation_ge_metrics", avro_schema()
157111
)
158112
feast_client.apply_entity(entity)
159113
feast_client.apply_feature_table(feature_table)

tests/e2e/utils/common.py

Lines changed: 51 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,51 @@
1+
from feast import (
2+
Client,
3+
Entity,
4+
Feature,
5+
FeatureTable,
6+
FileSource,
7+
KafkaSource,
8+
ValueType,
9+
)
10+
from feast.data_format import AvroFormat, ParquetFormat
11+
from feast.pyspark.abc import SparkJobStatus
12+
from feast.wait import wait_retry_backoff
13+
14+
15+
def create_schema(kafka_broker, topic_name, feature_table_name, avro_schema):
16+
entity = Entity(name="key", description="Key", value_type=ValueType.INT64)
17+
feature_table = FeatureTable(
18+
name=feature_table_name,
19+
entities=["key"],
20+
features=[Feature("num", ValueType.INT64), Feature("set", ValueType.STRING)],
21+
batch_source=FileSource(
22+
event_timestamp_column="event_timestamp",
23+
file_format=ParquetFormat(),
24+
file_url="/dev/null",
25+
),
26+
stream_source=KafkaSource(
27+
event_timestamp_column="event_timestamp",
28+
bootstrap_servers=kafka_broker,
29+
message_format=AvroFormat(avro_schema),
30+
topic=topic_name,
31+
),
32+
)
33+
return entity, feature_table
34+
35+
36+
def start_job(feast_client: Client, feature_table: FeatureTable, pytestconfig):
37+
if pytestconfig.getoption("scheduled_streaming_job"):
38+
return
39+
40+
job = feast_client.start_stream_to_online_ingestion(feature_table)
41+
wait_retry_backoff(
42+
lambda: (None, job.get_status() == SparkJobStatus.IN_PROGRESS), 240
43+
)
44+
return job
45+
46+
47+
def stop_job(job, feast_client: Client, feature_table: FeatureTable):
48+
if job:
49+
job.cancel()
50+
else:
51+
feast_client.delete_feature_table(feature_table.name)

0 commit comments

Comments
 (0)