|
@@ -71,11 +71,6 @@ func TestWindow(t *testing.T) {
|
|
|
}},
|
|
|
},
|
|
|
M: map[string]interface{}{
|
|
|
- "op_1_preprocessor_demo_0_exceptions_total": int64(0),
|
|
|
- "op_1_preprocessor_demo_0_process_latency_us": int64(0),
|
|
|
- "op_1_preprocessor_demo_0_records_in_total": int64(5),
|
|
|
- "op_1_preprocessor_demo_0_records_out_total": int64(5),
|
|
|
-
|
|
|
"op_3_project_0_exceptions_total": int64(0),
|
|
|
"op_3_project_0_process_latency_us": int64(0),
|
|
|
"op_3_project_0_records_in_total": int64(4),
|
|
@@ -111,11 +106,6 @@ func TestWindow(t *testing.T) {
|
|
|
}},
|
|
|
},
|
|
|
M: map[string]interface{}{
|
|
|
- "op_1_preprocessor_demo_0_exceptions_total": int64(0),
|
|
|
- "op_1_preprocessor_demo_0_process_latency_us": int64(0),
|
|
|
- "op_1_preprocessor_demo_0_records_in_total": int64(5),
|
|
|
- "op_1_preprocessor_demo_0_records_out_total": int64(5),
|
|
|
-
|
|
|
"op_4_project_0_exceptions_total": int64(0),
|
|
|
"op_4_project_0_process_latency_us": int64(0),
|
|
|
"op_4_project_0_records_in_total": int64(2),
|
|
@@ -206,16 +196,6 @@ func TestWindow(t *testing.T) {
|
|
|
}},
|
|
|
},
|
|
|
M: map[string]interface{}{
|
|
|
- "op_1_preprocessor_demo_0_exceptions_total": int64(0),
|
|
|
- "op_1_preprocessor_demo_0_process_latency_us": int64(0),
|
|
|
- "op_1_preprocessor_demo_0_records_in_total": int64(5),
|
|
|
- "op_1_preprocessor_demo_0_records_out_total": int64(5),
|
|
|
-
|
|
|
- "op_2_preprocessor_demo1_0_exceptions_total": int64(0),
|
|
|
- "op_2_preprocessor_demo1_0_process_latency_us": int64(0),
|
|
|
- "op_2_preprocessor_demo1_0_records_in_total": int64(5),
|
|
|
- "op_2_preprocessor_demo1_0_records_out_total": int64(5),
|
|
|
-
|
|
|
"op_5_project_0_exceptions_total": int64(0),
|
|
|
"op_5_project_0_process_latency_us": int64(0),
|
|
|
"op_5_project_0_records_in_total": int64(8),
|
|
@@ -246,13 +226,11 @@ func TestWindow(t *testing.T) {
|
|
|
T: &topo.PrintableTopo{
|
|
|
Sources: []string{"source_demo", "source_demo1"},
|
|
|
Edges: map[string][]string{
|
|
|
- "source_demo": {"op_1_preprocessor_demo"},
|
|
|
- "source_demo1": {"op_2_preprocessor_demo1"},
|
|
|
- "op_1_preprocessor_demo": {"op_3_window"},
|
|
|
- "op_2_preprocessor_demo1": {"op_3_window"},
|
|
|
- "op_3_window": {"op_4_join"},
|
|
|
- "op_4_join": {"op_5_project"},
|
|
|
- "op_5_project": {"sink_mockSink"},
|
|
|
+ "source_demo": {"op_3_window"},
|
|
|
+ "source_demo1": {"op_3_window"},
|
|
|
+ "op_3_window": {"op_4_join"},
|
|
|
+ "op_4_join": {"op_5_project"},
|
|
|
+ "op_5_project": {"sink_mockSink"},
|
|
|
},
|
|
|
},
|
|
|
}, {
|
|
@@ -292,11 +270,6 @@ func TestWindow(t *testing.T) {
|
|
|
}},
|
|
|
},
|
|
|
M: map[string]interface{}{
|
|
|
- "op_1_preprocessor_demo_0_exceptions_total": int64(0),
|
|
|
- "op_1_preprocessor_demo_0_process_latency_us": int64(0),
|
|
|
- "op_1_preprocessor_demo_0_records_in_total": int64(5),
|
|
|
- "op_1_preprocessor_demo_0_records_out_total": int64(5),
|
|
|
-
|
|
|
"op_5_project_0_exceptions_total": int64(0),
|
|
|
"op_5_project_0_process_latency_us": int64(0),
|
|
|
"op_5_project_0_records_in_total": int64(5),
|
|
@@ -348,11 +321,6 @@ func TestWindow(t *testing.T) {
|
|
|
}},
|
|
|
},
|
|
|
M: map[string]interface{}{
|
|
|
- "op_1_preprocessor_sessionDemo_0_exceptions_total": int64(0),
|
|
|
- "op_1_preprocessor_sessionDemo_0_process_latency_us": int64(0),
|
|
|
- "op_1_preprocessor_sessionDemo_0_records_in_total": int64(11),
|
|
|
- "op_1_preprocessor_sessionDemo_0_records_out_total": int64(11),
|
|
|
-
|
|
|
"op_3_project_0_exceptions_total": int64(0),
|
|
|
"op_3_project_0_process_latency_us": int64(0),
|
|
|
"op_3_project_0_records_in_total": int64(4),
|
|
@@ -418,16 +386,6 @@ func TestWindow(t *testing.T) {
|
|
|
}},
|
|
|
},
|
|
|
M: map[string]interface{}{
|
|
|
- "op_1_preprocessor_demo_0_exceptions_total": int64(0),
|
|
|
- "op_1_preprocessor_demo_0_process_latency_us": int64(0),
|
|
|
- "op_1_preprocessor_demo_0_records_in_total": int64(5),
|
|
|
- "op_1_preprocessor_demo_0_records_out_total": int64(5),
|
|
|
-
|
|
|
- "op_2_preprocessor_demo1_0_exceptions_total": int64(0),
|
|
|
- "op_2_preprocessor_demo1_0_process_latency_us": int64(0),
|
|
|
- "op_2_preprocessor_demo1_0_records_in_total": int64(5),
|
|
|
- "op_2_preprocessor_demo1_0_records_out_total": int64(5),
|
|
|
-
|
|
|
"op_5_project_0_exceptions_total": int64(0),
|
|
|
"op_5_project_0_process_latency_us": int64(0),
|
|
|
"op_5_project_0_records_in_total": int64(8),
|
|
@@ -489,11 +447,6 @@ func TestWindow(t *testing.T) {
|
|
|
}},
|
|
|
},
|
|
|
M: map[string]interface{}{
|
|
|
- "op_1_preprocessor_demoError_0_exceptions_total": int64(3),
|
|
|
- "op_1_preprocessor_demoError_0_process_latency_us": int64(0),
|
|
|
- "op_1_preprocessor_demoError_0_records_in_total": int64(5),
|
|
|
- "op_1_preprocessor_demoError_0_records_out_total": int64(2),
|
|
|
-
|
|
|
"op_3_project_0_exceptions_total": int64(3),
|
|
|
"op_3_project_0_process_latency_us": int64(0),
|
|
|
"op_3_project_0_records_in_total": int64(6),
|
|
@@ -503,7 +456,7 @@ func TestWindow(t *testing.T) {
|
|
|
"sink_mockSink_0_records_in_total": int64(6),
|
|
|
"sink_mockSink_0_records_out_total": int64(6),
|
|
|
|
|
|
- "source_demoError_0_exceptions_total": int64(0),
|
|
|
+ "source_demoError_0_exceptions_total": int64(3),
|
|
|
"source_demoError_0_records_in_total": int64(5),
|
|
|
"source_demoError_0_records_out_total": int64(5),
|
|
|
|
|
@@ -525,11 +478,6 @@ func TestWindow(t *testing.T) {
|
|
|
}},
|
|
|
},
|
|
|
M: map[string]interface{}{
|
|
|
- "op_1_preprocessor_demo_0_exceptions_total": int64(0),
|
|
|
- "op_1_preprocessor_demo_0_process_latency_us": int64(0),
|
|
|
- "op_1_preprocessor_demo_0_records_in_total": int64(5),
|
|
|
- "op_1_preprocessor_demo_0_records_out_total": int64(5),
|
|
|
-
|
|
|
"op_5_project_0_exceptions_total": int64(0),
|
|
|
"op_5_project_0_process_latency_us": int64(0),
|
|
|
"op_5_project_0_records_in_total": int64(1),
|
|
@@ -592,11 +540,6 @@ func TestWindow(t *testing.T) {
|
|
|
}},
|
|
|
},
|
|
|
M: map[string]interface{}{
|
|
|
- "op_1_preprocessor_demo_0_exceptions_total": int64(0),
|
|
|
- "op_1_preprocessor_demo_0_process_latency_us": int64(0),
|
|
|
- "op_1_preprocessor_demo_0_records_in_total": int64(5),
|
|
|
- "op_1_preprocessor_demo_0_records_out_total": int64(5),
|
|
|
-
|
|
|
"op_4_project_0_exceptions_total": int64(0),
|
|
|
"op_4_project_0_process_latency_us": int64(0),
|
|
|
"op_4_project_0_records_in_total": int64(4),
|
|
@@ -624,11 +567,6 @@ func TestWindow(t *testing.T) {
|
|
|
}},
|
|
|
},
|
|
|
M: map[string]interface{}{
|
|
|
- "op_1_preprocessor_demo_0_exceptions_total": int64(0),
|
|
|
- "op_1_preprocessor_demo_0_process_latency_us": int64(0),
|
|
|
- "op_1_preprocessor_demo_0_records_in_total": int64(5),
|
|
|
- "op_1_preprocessor_demo_0_records_out_total": int64(5),
|
|
|
-
|
|
|
"op_3_project_0_exceptions_total": int64(0),
|
|
|
"op_3_project_0_process_latency_us": int64(0),
|
|
|
"op_3_project_0_records_in_total": int64(1),
|
|
@@ -660,11 +598,6 @@ func TestWindow(t *testing.T) {
|
|
|
}}, {{}},
|
|
|
},
|
|
|
M: map[string]interface{}{
|
|
|
- "op_1_preprocessor_demo_0_exceptions_total": int64(0),
|
|
|
- "op_1_preprocessor_demo_0_process_latency_us": int64(0),
|
|
|
- "op_1_preprocessor_demo_0_records_in_total": int64(5),
|
|
|
- "op_1_preprocessor_demo_0_records_out_total": int64(5),
|
|
|
-
|
|
|
"op_3_project_0_exceptions_total": int64(0),
|
|
|
"op_3_project_0_process_latency_us": int64(0),
|
|
|
"op_3_project_0_records_in_total": int64(5),
|
|
@@ -710,14 +643,6 @@ func TestWindow(t *testing.T) {
|
|
|
"op_2_filter_0_records_in_total": int64(5),
|
|
|
"op_2_filter_0_records_out_total": int64(3),
|
|
|
|
|
|
- "op_1_preprocessor_demo_0_exceptions_total": int64(0),
|
|
|
- "op_1_preprocessor_demo_0_records_in_total": int64(5),
|
|
|
- "op_1_preprocessor_demo_0_records_out_total": int64(5),
|
|
|
-
|
|
|
- "op_4_tableprocessor_table1_0_exceptions_total": int64(0),
|
|
|
- "op_4_tableprocessor_table1_0_records_in_total": int64(4),
|
|
|
- "op_4_tableprocessor_table1_0_records_out_total": int64(1),
|
|
|
-
|
|
|
"op_5_filter_0_exceptions_total": int64(0),
|
|
|
"op_5_filter_0_records_in_total": int64(1),
|
|
|
"op_5_filter_0_records_out_total": int64(1),
|
|
@@ -743,7 +668,7 @@ func TestWindow(t *testing.T) {
|
|
|
|
|
|
"source_table1_0_exceptions_total": int64(0),
|
|
|
"source_table1_0_records_in_total": int64(4),
|
|
|
- "source_table1_0_records_out_total": int64(4),
|
|
|
+ "source_table1_0_records_out_total": int64(1),
|
|
|
},
|
|
|
},
|
|
|
}
|
|
@@ -815,11 +740,6 @@ func TestEventWindow(t *testing.T) {
|
|
|
}},
|
|
|
},
|
|
|
M: map[string]interface{}{
|
|
|
- "op_1_preprocessor_demoE_0_exceptions_total": int64(0),
|
|
|
- "op_1_preprocessor_demoE_0_process_latency_us": int64(0),
|
|
|
- "op_1_preprocessor_demoE_0_records_in_total": int64(6),
|
|
|
- "op_1_preprocessor_demoE_0_records_out_total": int64(6),
|
|
|
-
|
|
|
"op_3_project_0_exceptions_total": int64(0),
|
|
|
"op_3_project_0_process_latency_us": int64(0),
|
|
|
"op_3_project_0_records_in_total": int64(5),
|
|
@@ -856,11 +776,6 @@ func TestEventWindow(t *testing.T) {
|
|
|
}},
|
|
|
},
|
|
|
M: map[string]interface{}{
|
|
|
- "op_1_preprocessor_demoE_0_exceptions_total": int64(0),
|
|
|
- "op_1_preprocessor_demoE_0_process_latency_us": int64(0),
|
|
|
- "op_1_preprocessor_demoE_0_records_in_total": int64(6),
|
|
|
- "op_1_preprocessor_demoE_0_records_out_total": int64(6),
|
|
|
-
|
|
|
"op_4_project_0_exceptions_total": int64(0),
|
|
|
"op_4_project_0_process_latency_us": int64(0),
|
|
|
"op_4_project_0_records_in_total": int64(2),
|
|
@@ -919,16 +834,6 @@ func TestEventWindow(t *testing.T) {
|
|
|
}},
|
|
|
},
|
|
|
M: map[string]interface{}{
|
|
|
- "op_1_preprocessor_demoE_0_exceptions_total": int64(0),
|
|
|
- "op_1_preprocessor_demoE_0_process_latency_us": int64(0),
|
|
|
- "op_1_preprocessor_demoE_0_records_in_total": int64(6),
|
|
|
- "op_1_preprocessor_demoE_0_records_out_total": int64(6),
|
|
|
-
|
|
|
- "op_2_preprocessor_demo1E_0_exceptions_total": int64(0),
|
|
|
- "op_2_preprocessor_demo1E_0_process_latency_us": int64(0),
|
|
|
- "op_2_preprocessor_demo1E_0_records_in_total": int64(6),
|
|
|
- "op_2_preprocessor_demo1E_0_records_out_total": int64(6),
|
|
|
-
|
|
|
"op_5_project_0_exceptions_total": int64(0),
|
|
|
"op_5_project_0_process_latency_us": int64(0),
|
|
|
"op_5_project_0_records_in_total": int64(5),
|
|
@@ -995,11 +900,6 @@ func TestEventWindow(t *testing.T) {
|
|
|
}},
|
|
|
},
|
|
|
M: map[string]interface{}{
|
|
|
- "op_1_preprocessor_demoE_0_exceptions_total": int64(0),
|
|
|
- "op_1_preprocessor_demoE_0_process_latency_us": int64(0),
|
|
|
- "op_1_preprocessor_demoE_0_records_in_total": int64(6),
|
|
|
- "op_1_preprocessor_demoE_0_records_out_total": int64(6),
|
|
|
-
|
|
|
"op_5_project_0_exceptions_total": int64(0),
|
|
|
"op_5_project_0_process_latency_us": int64(0),
|
|
|
"op_5_project_0_records_in_total": int64(4),
|
|
@@ -1055,11 +955,6 @@ func TestEventWindow(t *testing.T) {
|
|
|
}},
|
|
|
},
|
|
|
M: map[string]interface{}{
|
|
|
- "op_1_preprocessor_sessionDemoE_0_exceptions_total": int64(0),
|
|
|
- "op_1_preprocessor_sessionDemoE_0_process_latency_us": int64(0),
|
|
|
- "op_1_preprocessor_sessionDemoE_0_records_in_total": int64(12),
|
|
|
- "op_1_preprocessor_sessionDemoE_0_records_out_total": int64(12),
|
|
|
-
|
|
|
"op_3_project_0_exceptions_total": int64(0),
|
|
|
"op_3_project_0_process_latency_us": int64(0),
|
|
|
"op_3_project_0_records_in_total": int64(4),
|
|
@@ -1100,16 +995,6 @@ func TestEventWindow(t *testing.T) {
|
|
|
}},
|
|
|
},
|
|
|
M: map[string]interface{}{
|
|
|
- "op_1_preprocessor_demoE_0_exceptions_total": int64(0),
|
|
|
- "op_1_preprocessor_demoE_0_process_latency_us": int64(0),
|
|
|
- "op_1_preprocessor_demoE_0_records_in_total": int64(6),
|
|
|
- "op_1_preprocessor_demoE_0_records_out_total": int64(6),
|
|
|
-
|
|
|
- "op_2_preprocessor_demo1E_0_exceptions_total": int64(0),
|
|
|
- "op_2_preprocessor_demo1E_0_process_latency_us": int64(0),
|
|
|
- "op_2_preprocessor_demo1E_0_records_in_total": int64(6),
|
|
|
- "op_2_preprocessor_demo1E_0_records_out_total": int64(6),
|
|
|
-
|
|
|
"op_5_project_0_exceptions_total": int64(0),
|
|
|
"op_5_project_0_process_latency_us": int64(0),
|
|
|
"op_5_project_0_records_in_total": int64(5),
|
|
@@ -1172,11 +1057,6 @@ func TestEventWindow(t *testing.T) {
|
|
|
}},
|
|
|
},
|
|
|
M: map[string]interface{}{
|
|
|
- "op_1_preprocessor_demoErr_0_exceptions_total": int64(1),
|
|
|
- "op_1_preprocessor_demoErr_0_process_latency_us": int64(0),
|
|
|
- "op_1_preprocessor_demoErr_0_records_in_total": int64(6),
|
|
|
- "op_1_preprocessor_demoErr_0_records_out_total": int64(5),
|
|
|
-
|
|
|
"op_3_project_0_exceptions_total": int64(1),
|
|
|
"op_3_project_0_process_latency_us": int64(0),
|
|
|
"op_3_project_0_records_in_total": int64(6),
|
|
@@ -1186,7 +1066,7 @@ func TestEventWindow(t *testing.T) {
|
|
|
"sink_mockSink_0_records_in_total": int64(6),
|
|
|
"sink_mockSink_0_records_out_total": int64(6),
|
|
|
|
|
|
- "source_demoErr_0_exceptions_total": int64(0),
|
|
|
+ "source_demoErr_0_exceptions_total": int64(1),
|
|
|
"source_demoErr_0_records_in_total": int64(6),
|
|
|
"source_demoErr_0_records_out_total": int64(6),
|
|
|
|
|
@@ -1242,11 +1122,6 @@ func TestEventWindow(t *testing.T) {
|
|
|
}},
|
|
|
},
|
|
|
M: map[string]interface{}{
|
|
|
- "op_1_preprocessor_sessionDemoE_0_exceptions_total": int64(0),
|
|
|
- "op_1_preprocessor_sessionDemoE_0_process_latency_us": int64(0),
|
|
|
- "op_1_preprocessor_sessionDemoE_0_records_in_total": int64(12),
|
|
|
- "op_1_preprocessor_sessionDemoE_0_records_out_total": int64(12),
|
|
|
-
|
|
|
"op_3_project_0_exceptions_total": int64(0),
|
|
|
"op_3_project_0_process_latency_us": int64(0),
|
|
|
"op_3_project_0_records_in_total": int64(4),
|
|
@@ -1306,11 +1181,6 @@ func TestEventWindow(t *testing.T) {
|
|
|
}},
|
|
|
},
|
|
|
M: map[string]interface{}{
|
|
|
- "op_1_preprocessor_demoE_0_exceptions_total": int64(0),
|
|
|
- "op_1_preprocessor_demoE_0_process_latency_us": int64(0),
|
|
|
- "op_1_preprocessor_demoE_0_records_in_total": int64(6),
|
|
|
- "op_1_preprocessor_demoE_0_records_out_total": int64(6),
|
|
|
-
|
|
|
"op_3_project_0_exceptions_total": int64(0),
|
|
|
"op_3_project_0_process_latency_us": int64(0),
|
|
|
"op_3_project_0_records_in_total": int64(5),
|
|
@@ -1375,11 +1245,6 @@ func TestWindowError(t *testing.T) {
|
|
|
}, {}},
|
|
|
},
|
|
|
M: map[string]interface{}{
|
|
|
- "op_1_preprocessor_ldemo_0_exceptions_total": int64(0),
|
|
|
- "op_1_preprocessor_ldemo_0_process_latency_us": int64(0),
|
|
|
- "op_1_preprocessor_ldemo_0_records_in_total": int64(5),
|
|
|
- "op_1_preprocessor_ldemo_0_records_out_total": int64(5),
|
|
|
-
|
|
|
"op_3_project_0_exceptions_total": int64(1),
|
|
|
"op_3_project_0_process_latency_us": int64(0),
|
|
|
"op_3_project_0_records_in_total": int64(2),
|
|
@@ -1412,11 +1277,6 @@ func TestWindowError(t *testing.T) {
|
|
|
}},
|
|
|
},
|
|
|
M: map[string]interface{}{
|
|
|
- "op_1_preprocessor_ldemo_0_exceptions_total": int64(0),
|
|
|
- "op_1_preprocessor_ldemo_0_process_latency_us": int64(0),
|
|
|
- "op_1_preprocessor_ldemo_0_records_in_total": int64(5),
|
|
|
- "op_1_preprocessor_ldemo_0_records_out_total": int64(5),
|
|
|
-
|
|
|
"op_4_project_0_exceptions_total": int64(1),
|
|
|
"op_4_project_0_process_latency_us": int64(0),
|
|
|
"op_4_project_0_records_in_total": int64(3),
|
|
@@ -1471,16 +1331,6 @@ func TestWindowError(t *testing.T) {
|
|
|
}},
|
|
|
},
|
|
|
M: map[string]interface{}{
|
|
|
- "op_1_preprocessor_ldemo_0_exceptions_total": int64(0),
|
|
|
- "op_1_preprocessor_ldemo_0_process_latency_us": int64(0),
|
|
|
- "op_1_preprocessor_ldemo_0_records_in_total": int64(5),
|
|
|
- "op_1_preprocessor_ldemo_0_records_out_total": int64(5),
|
|
|
-
|
|
|
- "op_2_preprocessor_ldemo1_0_exceptions_total": int64(0),
|
|
|
- "op_2_preprocessor_ldemo1_0_process_latency_us": int64(0),
|
|
|
- "op_2_preprocessor_ldemo1_0_records_in_total": int64(5),
|
|
|
- "op_2_preprocessor_ldemo1_0_records_out_total": int64(5),
|
|
|
-
|
|
|
"op_5_project_0_exceptions_total": int64(3),
|
|
|
"op_5_project_0_process_latency_us": int64(0),
|
|
|
"op_5_project_0_records_in_total": int64(8),
|
|
@@ -1525,11 +1375,6 @@ func TestWindowError(t *testing.T) {
|
|
|
}, {}},
|
|
|
},
|
|
|
M: map[string]interface{}{
|
|
|
- "op_1_preprocessor_ldemo_0_exceptions_total": int64(0),
|
|
|
- "op_1_preprocessor_ldemo_0_process_latency_us": int64(0),
|
|
|
- "op_1_preprocessor_ldemo_0_records_in_total": int64(5),
|
|
|
- "op_1_preprocessor_ldemo_0_records_out_total": int64(5),
|
|
|
-
|
|
|
"op_6_project_0_exceptions_total": int64(3),
|
|
|
"op_6_project_0_process_latency_us": int64(0),
|
|
|
"op_6_project_0_records_in_total": int64(5),
|
|
@@ -1574,11 +1419,6 @@ func TestWindowError(t *testing.T) {
|
|
|
}},
|
|
|
},
|
|
|
M: map[string]interface{}{
|
|
|
- "op_1_preprocessor_ldemo_0_exceptions_total": int64(0),
|
|
|
- "op_1_preprocessor_ldemo_0_process_latency_us": int64(0),
|
|
|
- "op_1_preprocessor_ldemo_0_records_in_total": int64(5),
|
|
|
- "op_1_preprocessor_ldemo_0_records_out_total": int64(5),
|
|
|
-
|
|
|
"op_4_project_0_exceptions_total": int64(1),
|
|
|
"op_4_project_0_process_latency_us": int64(0),
|
|
|
"op_4_project_0_records_in_total": int64(4),
|