- 07 1月, 2021 12 次提交
-
-
由 Aljoscha Krettek 提交于
This adds an ITCase because we need to check that all the components work together and that the wiring works correctly. This uses the previously added funtionality to specify that given inputs should be processed before other inputs and should not be sorted.
-
由 Aljoscha Krettek 提交于
This doesn't change the actual behavior, we still set the same "sorted" setting on both inputs. We will add tests and actually change the behavior in a follow-up commit.
-
由 Aljoscha Krettek 提交于
We allow the operation but it is a no-op because there never is any keyed state because we process all broadcast input before starting to process the keyed/other side.
-
由 Aljoscha Krettek 提交于
This will allow processing the broadcast side of a broadcast operator first, before processing the keyed side that requires sorting for stateful BATCH execution. For now, the wiring from the API is not there, this will be added in follow-up changes.
-
由 Aljoscha Krettek 提交于
Note: Broadcast operations in BATCH mode don't yet work with this change. This needs follow-up changes from later commits. We just lay the groundwork here and keep the same functionality. Before, we were creating operators in BroadcastConnectedStreams eagerly. Now, the transformation holds the user function and we add a Translator that creates the "physical" operators when translating the graph of Transformations. We do this so that we can translate differently based on whether we're in BATCH or STREAMING mode.
-
由 Aljoscha Krettek 提交于
Because that's what it does.
-
由 Dian Fu 提交于
-
由 龙三 提交于
[FLINK-20766][table-planner-blink] Introduce BatchPhysicalSortLimit, and make BatchExecSortLimit only extended from ExecNode This closes #14502
-
由 龙三 提交于
[FLINK-20766][table-planner-blink] Introduce BatchPhysicalSort, and make BatchExeSort only extended from ExecNode This closes #14502
-
由 龙三 提交于
[FLINK-20766][table-planner-blink] Introduce StreamPhysicalTemporalSort, and make StreamExecTemporalSort only extended from ExecNode. This closes #14502
-
由 龙三 提交于
[FLINK-20766][table-planner-blink] Introduce StreamPhysicalSortLimit, and make StreamExecSortLimit only extends from ExecNode. This closes #14502
-
由 龙三 提交于
[FLINK-20766][table-planner-blink] Introduce StreamPhysicalSort, and make StreamExecSort only extended from ExecNode This closes #14502
-
- 06 1月, 2021 18 次提交
-
-
由 龙三 提交于
[FLINK-20782][table-planner-blink] Introduce BatchPhysicalRank, and make BatchExecRank only extended from ExecNode This closes #14506
-
由 Piotr Nowojski 提交于
In particularly, if task is idling forever, as there are no new records incomming previous version would report idleTime as 0% and busyTime as 100% which is incorrect. In this version, idleTime metric is aware that idling period has started and can take that into account when updating it's value.
-
由 Piotr Nowojski 提交于
-
由 Piotr Nowojski 提交于
-
由 Piotr Nowojski 提交于
It's defined as inverted value of idleTimeMsPerSecond
-
由 Piotr Nowojski 提交于
-
由 Piotr Nowojski 提交于
-
由 Piotr Nowojski 提交于
-
由 Jie Wang 提交于
This closes #14467
-
由 Jark Wu 提交于
This closes #14563
-
由 godfreyhe 提交于
[FLINK-20737][table-planner-blink] Introduce StreamPhysicalPythonGroupTableAggregate, and make StreamExecPythonGroupTableAggregate only extended from ExecNode This closes #14478
-
由 godfreyhe 提交于
[FLINK-20737][table-planner-blink] Introduce StreamPhysicalGroupTableAggregate, and make StreamExecGroupTableAggregate only extended from ExecNode This closes #14478
-
由 godfreyhe 提交于
[FLINK-20737][table-planner-blink] Introduce StreamPhysicalIncrementalGroupAggregate, and make StreamExecIncrementalGroupAggregate only extended from ExecNode This closes #14478
-
由 godfreyhe 提交于
[FLINK-20737][table-planner-blink] Introduce StreamPhysicalGlobalGroupAggregate, and make StreamExecGlobalGroupAggregate only extended from ExecNode This closes #14478
-
由 godfreyhe 提交于
[FLINK-20737][table-planner-blink] Introduce StreamPhysicalLocalGroupAggregate, and make StreamExecLocalGroupAggregate only extended from ExecNode This closes #14478
-
由 godfreyhe 提交于
[FLINK-20737][table-planner-blink] Introduce StreamPhysicalPythonGroupAggregate, and make StreamExecPythonGroupAggregate only extended from ExecNode This closes #14478
-
由 godfreyhe 提交于
[FLINK-20737][table-planner-blink] Introduce StreamPhysicalGroupAggregate, and make StreamExecGroupAggregate only extended from ExecNode This closes #14478
-
由 godfreyhe 提交于
This closes #14478
-
- 05 1月, 2021 10 次提交
-
-
由 Till Rohrmann 提交于
Removing the ProgrammedSlotProvider since it is no longer used.
-
由 Till Rohrmann 提交于
Moved the factories of the CompletedCheckpointStore and the CheckpointIDCounter to SchedulerUtils in order to make them reusable.
-
由 Till Rohrmann 提交于
-
由 Till Rohrmann 提交于
By moving the shut down of checkpoint services out of the CheckpointCoordinator, it is now possible to reuse these services across different CheckpointCoordinators. This closes #14553.
-
由 Till Rohrmann 提交于
-
由 Till Rohrmann 提交于
-
由 Till Rohrmann 提交于
-
由 Till Rohrmann 提交于
[hotfix][tests] Replace explicit ExecutionGraphBuilder.buildGraph calls with TestingExecutionGraphBuilder
-
由 Till Rohrmann 提交于
-
由 Till Rohrmann 提交于
[FLINK-20846] Factor out CompletedCheckpointStore and CheckpointIDCounter creation in ExecutionGraphBuilder.buildGraph
-