Skip to content

Commit 1ec8bad

Browse files
committed
[Timeseries] TS-36-Remove-tftlayer-post-process
1 parent 2b9b40e commit 1ec8bad

6 files changed

Lines changed: 96 additions & 84 deletions

File tree

‎timeseries-streaming/timeseries-python-applications/README.MD‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -75,7 +75,7 @@ That job will generate 86400 data points, which will be aggregated to 43200 exam
7575
Once that is done, change the information in the [config.py](ml_pipeline_examples/sin_wave_example/config.py) to match your local env.
7676
Run the command with the virtual-env activated:
7777
```
78-
python ml_pipeline_examples/sin_wave_example/timeseries_local_simple_data.py
78+
python ml_pipeline_examples/sin_wave_example/training/timeseries_local_sin_wave.py
7979
```
8080

8181
This will output a serving_model_dir and a tf_transform_graph under the location you specified for ```PIPELINE_ROOT``` in the config.py file. With this you can now follow the rest of the steps outlines in Option 1 but using your own model.

‎timeseries-streaming/timeseries-python-applications/ml_pipeline/timeseries/encoder_decoder/encoder_decoder_run_fn.py‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,7 @@
1616
from typing import Text
1717
import tensorflow as tf
1818
import tensorflow_transform as tft
19-
import timeseries.encoder_decoder.encoder_decoder_model as encoder_decoder_model
19+
import ml_pipeline.timeseries.encoder_decoder.encoder_decoder_model as encoder_decoder_model
2020

2121
from tfx.components.trainer.executor import TrainerFnArgs
2222

@@ -69,7 +69,7 @@ def serve_tf_examples_fn(serialized_tf_examples):
6969
serialized_tf_examples, feature_spec)
7070
transformed_features = create_training_data(
7171
model.tft_layer(parsed_features))
72-
return model(transformed_features)
72+
return model(transformed_features), transformed_features
7373

7474
return serve_tf_examples_fn
7575

‎timeseries-streaming/timeseries-python-applications/ml_pipeline/timeseries/encoder_decoder/transforms/process_encdec_inf_rtn.py‎

Lines changed: 22 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -42,8 +42,9 @@ def __init__(self, config: Dict[Text, Any], batching_size: int = 1000):
4242
self.batching_size = batching_size
4343

4444
def setup(self):
45+
# TODO switch to shared.py
4546
self.transform_output = tft.TFTransformOutput(self.tf_transform_graph_dir)
46-
self.tft_layer = self.transform_output.transform_features_layer()
47+
# self.tft_layer = self.transform_output.transform_features_layer()
4748

4849
def start_bundle(self):
4950
self.batch: [WindowedValue] = []
@@ -82,15 +83,17 @@ def process_result(self, element: [WindowedValue]):
8283
element_value = [k.value for k in element]
8384
processed_inputs = []
8485
request_inputs = []
86+
request_inputs_tft = []
8587
request_outputs = []
8688

8789
for k in element_value:
8890
request_inputs.append(
8991
k.predict_log.request.inputs['examples'].string_val[0])
9092
request_outputs.append(k.predict_log.response.outputs['output_0'])
93+
request_inputs_tft.append(k.predict_log.response.outputs['output_1'])
9194

9295
# The output of tf.io.parse_example is a set of feature tensors which
93-
# have shape for non Metadata of [batch,
96+
# have shape for non-metadata values of [batch,
9497
# timestep]
9598

9699
batched_example = tf.io.parse_example(
@@ -99,7 +102,9 @@ def process_result(self, element: [WindowedValue]):
99102
# The tft layer gives us two labels 'FLOAT32' and 'LABEL' which have
100103
# shape [batch, timestep, model_features]
101104

102-
inputs = self.tft_layer(batched_example)
105+
# inputs = self.tft_layer(batched_example)
106+
# TODO make work with batches
107+
inputs = request_inputs_tft
103108

104109
# Determine which of the features was used in the model
105110
feature_labels = timeseries_transform_utils.create_feature_list_from_list(
@@ -115,7 +120,9 @@ def process_result(self, element: [WindowedValue]):
115120
batched_example['METADATA_SPAN_END_TS']).numpy()
116121

117122
batch_pos = 0
118-
for batch_input in inputs['LABEL'].numpy():
123+
for batch_input in inputs:
124+
# TODO Currently only supports model runinf of batch size 1
125+
batch_input = tf.make_ndarray(batch_input)[0]
119126
# Get the Metadata from the original request
120127
span_start_timestamp = datetime.fromtimestamp(
121128
metadata_span_start_timestamp[batch_pos][0] / 1000)
@@ -150,9 +157,9 @@ def process_result(self, element: [WindowedValue]):
150157
label = (feature_labels[model_feature_pos])
151158

152159
# The num of features should == number of results
153-
if len(feature_labels) != len(last_timestep_input):
160+
if len(feature_labels) != len(last_timestep_output):
154161
raise ValueError(f'Features list {feature_labels} in config is '
155-
f'len {len(feature_labels)} which '
162+
f'length {len(feature_labels)} which '
156163
f'does not match output length '
157164
f'{len(last_timestep_output)} '
158165
f' This normally is a result of using a configuration '
@@ -198,8 +205,9 @@ class CheckAnomalous(beam.DoFn):
198205
"""
199206

200207
# TODO(BEAM-6158): Revert the workaround once we can pickle super() on py3.
201-
def __init__(self, threshold: float = 0.05):
208+
def __init__(self, output_false: bool = True, threshold: float = 0.05):
202209
beam.DoFn.__init__(self)
210+
self.output_false = output_false
203211
self.threshold = threshold
204212

205213
def process(self, element: Dict[Text, Any], *unused_args, **unused_kwargs):
@@ -208,12 +216,17 @@ def process(self, element: Dict[Text, Any], *unused_args, **unused_kwargs):
208216
'span_end_timestamp': element['span_end_timestamp']
209217
}
210218

219+
anomaly_found = False
220+
211221
for key, value in element['feature_results'].items():
212222
input_value = value['input_value']
213223
output_value = value['output_value']
214224
diff = abs(input_value - output_value)
215225
value.update({'diff': diff})
216226
if not key.endswith('-TIMESTAMP'):
217-
value.update({'anomaly': diff > self.threshold})
227+
if diff > self.threshold:
228+
value.update({'anomaly': True})
229+
anomaly_found = True
218230
result.update({key: value})
219-
yield result
231+
if anomaly_found:
232+
yield result

‎timeseries-streaming/timeseries-python-applications/ml_pipeline/timeseries/utils/timeseries_transform_utils_test.py‎

Lines changed: 66 additions & 67 deletions
Original file line numberDiff line numberDiff line change
@@ -18,18 +18,17 @@
1818
import numpy
1919
from datetime import datetime
2020
import tensorflow as tf
21-
import timeseries.encoder_decoder.encoder_decoder_preprocessing as encoder_decoder_preprocessing
21+
import ml_pipeline.timeseries.encoder_decoder.encoder_decoder_preprocessing as encoder_decoder_preprocessing
2222
from scipy.stats import stats
2323
from tensorflow_transform.beam import tft_unit
2424
from tensorflow_transform.beam import impl as beam_impl
25-
from timeseries.utils import timeseries_transform_utils as ts_utils
26-
25+
from ml_pipeline.timeseries.utils import timeseries_transform_utils as ts_utils
2726

2827

2928
class BeamImplTest(tft_unit.TransformTestCase):
3029
def setUp(self):
3130
tf.compat.v1.logging.info(
32-
'Starting test case: %s', self._testMethodName)
31+
'Starting test case: %s', self._testMethodName)
3332

3433
self._context = beam_impl.Context(use_deep_copy_optimization=True)
3534
self._context.__enter__()
@@ -43,15 +42,15 @@ def _SkipIfExternalEnvironmentAnd(self, predicate, reason):
4342

4443
def testBasicType(self):
4544
config = {
46-
'timesteps': 3,
47-
'time_features': [],
48-
'features': ['a'],
49-
'enable_timestamp_features': False
45+
'timesteps': 3,
46+
'time_features': [],
47+
'features': ['a'],
48+
'enable_timestamp_features': False
5049
}
5150

5251
input_data = [{'a': [1000.0, 2000.0, 3000.0]}]
5352
input_metadata = tft_unit.metadata_from_feature_spec(
54-
{'a': tf.io.VarLenFeature(tf.float32)})
53+
{'a': tf.io.VarLenFeature(tf.float32)})
5554

5655
output = [[1000], [2000], [3000]]
5756

@@ -60,37 +59,37 @@ def testBasicType(self):
6059
expected_data = [{'Float32': output, 'LABEL': output}]
6160

6261
expected_metadata = tft_unit.metadata_from_feature_spec({
63-
'Float32': tf.io.FixedLenFeature([config['timesteps'], 1],
64-
tf.float32),
65-
'LABEL': tf.io.FixedLenFeature([config['timesteps'], 1],
66-
tf.float32)
62+
'Float32': tf.io.FixedLenFeature([config['timesteps'], 1],
63+
tf.float32),
64+
'LABEL': tf.io.FixedLenFeature([config['timesteps'], 1],
65+
tf.float32)
6766
})
6867

6968
preprocessing_fn = functools.partial(
70-
encoder_decoder_preprocessing.preprocessing_fn,
71-
custom_config=config)
69+
encoder_decoder_preprocessing.preprocessing_fn,
70+
custom_config=config)
7271

7372
self.assertAnalyzeAndTransformResults(
74-
input_data,
75-
input_metadata,
76-
preprocessing_fn,
77-
expected_data,
78-
expected_metadata)
73+
input_data,
74+
input_metadata,
75+
preprocessing_fn,
76+
expected_data,
77+
expected_metadata)
7978

8079
def testMixedType(self):
8180
config = {
82-
'timesteps': 3,
83-
'time_features': ['MINUTE', 'MONTH', 'HOUR', 'DAY', 'YEAR'],
84-
'features': ['a', 'b'],
85-
'enable_timestamp_features': False
81+
'timesteps': 3,
82+
'time_features': ['MINUTE', 'MONTH', 'HOUR', 'DAY', 'YEAR'],
83+
'features': ['a', 'b'],
84+
'enable_timestamp_features': False
8685
}
8786

8887
input_data = [{
89-
'a': [1000.0, 2000.0, 3000.0], 'b': [3000, 2000, 1000]
88+
'a': [1000.0, 2000.0, 3000.0], 'b': [3000, 2000, 1000]
9089
}]
9190
input_metadata = tft_unit.metadata_from_feature_spec({
92-
'a': tf.io.VarLenFeature(tf.float32),
93-
'b': tf.io.VarLenFeature(tf.int64)
91+
'a': tf.io.VarLenFeature(tf.float32),
92+
'b': tf.io.VarLenFeature(tf.int64)
9493
})
9594

9695
output = [[1000.0, 3000.0], [2000.0, 2000.0], [3000.0, 1000.0]]
@@ -100,43 +99,43 @@ def testMixedType(self):
10099
expected_data = [{'Float32': output, 'LABEL': output}]
101100

102101
expected_metadata = tft_unit.metadata_from_feature_spec({
103-
'Float32': tf.io.FixedLenFeature([config['timesteps'], 2],
104-
tf.float32),
105-
'LABEL': tf.io.FixedLenFeature([config['timesteps'], 2],
106-
tf.float32)
102+
'Float32': tf.io.FixedLenFeature([config['timesteps'], 2],
103+
tf.float32),
104+
'LABEL': tf.io.FixedLenFeature([config['timesteps'], 2],
105+
tf.float32)
107106
})
108107

109108
preprocessing_fn = functools.partial(
110-
encoder_decoder_preprocessing.preprocessing_fn,
111-
custom_config=config)
109+
encoder_decoder_preprocessing.preprocessing_fn,
110+
custom_config=config)
112111

113112
self.assertAnalyzeAndTransformResults(
114-
input_data,
115-
input_metadata,
116-
preprocessing_fn,
117-
expected_data,
118-
expected_metadata)
113+
input_data,
114+
input_metadata,
115+
preprocessing_fn,
116+
expected_data,
117+
expected_metadata)
119118

120119
def testWithTimeStamps(self):
121120

122121
config = {
123-
'timesteps': 2,
124-
'time_features': ['MINUTE', 'MONTH', 'HOUR', 'DAY', 'YEAR'],
125-
'features': ['float32', 'foo_TIMESTAMP'],
126-
'enable_timestamp_features': True
122+
'timesteps': 2,
123+
'time_features': ['MINUTE', 'MONTH', 'HOUR', 'DAY', 'YEAR'],
124+
'features': ['float32', 'foo_TIMESTAMP'],
125+
'enable_timestamp_features': True
127126
}
128127

129128
# The values will need to be different enough for the zscore not to nan
130129
timestamp_1 = int(datetime(2000, 1, 1, 0, 0, 0).timestamp())
131130
timestamp_2 = int(datetime(2001, 6, 15, 12, 30, 30).timestamp())
132131

133132
input_data = [{
134-
'float32': [1000.0, 2000.0],
135-
'foo_TIMESTAMP': [timestamp_1 * 1000, timestamp_2 * 1000]
133+
'float32': [1000.0, 2000.0],
134+
'foo_TIMESTAMP': [timestamp_1 * 1000, timestamp_2 * 1000]
136135
}]
137136
input_metadata = tft_unit.metadata_from_feature_spec({
138-
'float32': tf.io.VarLenFeature(tf.float32),
139-
'foo_TIMESTAMP': tf.io.VarLenFeature(tf.int64)
137+
'float32': tf.io.VarLenFeature(tf.float32),
138+
'foo_TIMESTAMP': tf.io.VarLenFeature(tf.int64)
140139
})
141140

142141
output_timestep_1 = self.create_transform_output(timestamp_1)
@@ -160,22 +159,22 @@ def testWithTimeStamps(self):
160159
expected_data = [{'Float32': output, 'LABEL': output}]
161160

162161
expected_metadata = tft_unit.metadata_from_feature_spec({
163-
'Float32': tf.io.FixedLenFeature([config['timesteps'], 11],
164-
tf.float32),
165-
'LABEL': tf.io.FixedLenFeature([config['timesteps'], 11],
166-
tf.float32)
162+
'Float32': tf.io.FixedLenFeature([config['timesteps'], 11],
163+
tf.float32),
164+
'LABEL': tf.io.FixedLenFeature([config['timesteps'], 11],
165+
tf.float32)
167166
})
168167

169168
preprocessing_fn = functools.partial(
170-
encoder_decoder_preprocessing.preprocessing_fn,
171-
custom_config=config)
169+
encoder_decoder_preprocessing.preprocessing_fn,
170+
custom_config=config)
172171

173172
self.assertAnalyzeAndTransformResults(
174-
input_data,
175-
input_metadata,
176-
preprocessing_fn,
177-
expected_data,
178-
expected_metadata)
173+
input_data,
174+
input_metadata,
175+
preprocessing_fn,
176+
expected_data,
177+
expected_metadata)
179178

180179
def create_transform_output(self, timestamp: int) -> [float]:
181180
# Needs to be in lexical order
@@ -195,16 +194,16 @@ def create_transform_output(self, timestamp: int) -> [float]:
195194
cos_year = math.cos(timestamp * (2.0 * math.pi / ts_utils.YEAR))
196195

197196
return [
198-
cos_day,
199-
cos_hour,
200-
cos_min,
201-
cos_month,
202-
cos_year,
203-
sin_day,
204-
sin_hour,
205-
sin_min,
206-
sin_month,
207-
sin_year
197+
cos_day,
198+
cos_hour,
199+
cos_min,
200+
cos_month,
201+
cos_year,
202+
sin_day,
203+
sin_hour,
204+
sin_min,
205+
sin_month,
206+
sin_year
208207
]
209208

210209

‎timeseries-streaming/timeseries-python-applications/ml_pipeline_examples/sin_wave_example/inference/stream_inference.py‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -60,7 +60,7 @@ def run(args, pipeline_args):
6060
'model_config': config.MODEL_CONFIG
6161
}))
6262
| beam.ParDo(
63-
process_encdec_inf_rtn.CheckAnomalous(threshold=0.7))
63+
process_encdec_inf_rtn.CheckAnomalous(output_false=False, threshold=0.7))
6464
| beam.ParDo(print))
6565

6666

‎timeseries-streaming/timeseries-python-applications/ml_pipeline_examples/sin_wave_example/training/timeseries_local_sin_wave.py‎

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -19,12 +19,12 @@
1919
from tfx.orchestration.beam.beam_dag_runner import BeamDagRunner
2020
from tfx.proto import trainer_pb2
2121

22-
import timeseries.pipeline_templates.timeseries_pipeline as pipeline
22+
import ml_pipeline.timeseries.pipeline_templates.timeseries_pipeline as pipeline
2323
from tfx.orchestration import metadata
2424
from absl import logging
2525

2626
import ml_pipeline_examples.sin_wave_example.config as config
27-
from timeseries.utils import timeseries_transform_utils
27+
from ml_pipeline.timeseries.utils import timeseries_transform_utils
2828

2929

3030
def run():
@@ -39,9 +39,9 @@ def run():
3939
tfx_pipeline = pipeline.create_pipeline(
4040
pipeline_name=config.PIPELINE_NAME,
4141
enable_cache=False,
42-
run_fn='timeseries.encoder_decoder.encoder_decoder_run_fn.run_fn',
42+
run_fn='ml_pipeline.timeseries.encoder_decoder.encoder_decoder_run_fn.run_fn',
4343
preprocessing_fn=
44-
'timeseries.encoder_decoder.encoder_decoder_preprocessing.preprocessing_fn',
44+
'ml_pipeline.timeseries.encoder_decoder.encoder_decoder_preprocessing.preprocessing_fn',
4545
data_path=config.SYNTHETIC_DATASET['local-raw'],
4646
pipeline_root=config.LOCAL_PIPELINE_ROOT,
4747
serving_model_dir=join(config.LOCAL_PIPELINE_ROOT, os.pathsep),

0 commit comments

Comments
 (0)