- 12 5月, 2020 30 次提交
-
-
由 Roman Khachatryan 提交于
Motivation: current behavior of forcing savepoints causes excess of max-concurrent-checkpoints limit. Which violates current Unaligned Checkpoints (UC) assumption of a single concurrent barrier. Changes: 1. Remove periodicTriggeringSuspended. Instead, if periodic request can't be executed now is ignored. 2. Check queue before executing a new request (not after). 3. Execute queued request when pending checkpoint is completed (instead of resuming timer). 4. Don't set CheckpointRequest.force in UC mode. 5. Change requests queue to PriorityQueue to prioritize savepoints. Limit its size. 6. For savepoints, if running in UC mode and limit is reached then enqueue request instead of forcing it 7. Don't throw trigger exceptions. CheckpointFailureManager ignores them
-
由 Roman Khachatryan 提交于
Pre-requsite refactoring to change decision logic
-
由 Roman Khachatryan 提交于
-
由 Jörn Kottmann 提交于
This also adds a test case to ensure this works as expected.
-
由 Guru Prasad 提交于
-
由 Zhenghua Gao 提交于
This closes #11892
-
由 Dawid Wysakowicz 提交于
This closes #11981
-
由 Dawid Wysakowicz 提交于
-
由 Aljoscha Krettek 提交于
Co-authored-by: NDawid Wysakowicz <dwysakowicz@apache.org>
-
由 Dawid Wysakowicz 提交于
parser
-
由 Jingsong Lee 提交于
This closes #12091
-
由 Leonard Xu 提交于
-
由 Dawid Wysakowicz 提交于
This PR adds a way to emit multiple records from KafkaDeserializationSchema. This is possible through a collector, which will buffer deserialized records in a queue and then emit all records atomically. The queue is reused for all incoming Kafka records to minimize creating new objects on the hot path.
-
由 Dawid Wysakowicz 提交于
-
由 Robert Metzger 提交于
This closes #12016
-
由 Yangze Guo 提交于
[FLINK-17536][core] Change the config option for slot max limitation to slotmanager.number-of-slots.max This closes #12067.
-
由 JingsongLi 提交于
-
由 Jingsong Lee 提交于
This closes #12004
-
由 yuzhao.cyz 提交于
[FLINK-17577][table-common] SinkFormat#createSinkFormat should use DynamicTableSink.Context as parameter This closes #12039
-
由 Shuiqiang Chen 提交于
[FLINK-17609][python] Execute the script directly when user specified the entry script with "-py" (#12079)
-
由 vthinkxie 提交于
This closes #11731.
-
由 Zhijiang 提交于
[FLINK-16536][network][checkpointing] Implement InputChannel state recovery for unaligned checkpoint During recovery process for unaligned checkpoint, the input channel state should also be recovered besides with existing operator states. We considered three guarantees during the implementation: 1. Make input recovery happen after the output recovery for providing more floating buffers on output side firstly. 2. Make partition request happen after input recovery for avoiding new data overtaking the previous state data. 3. Introduce a dedicated single IO executor for unspilling the channel state one by one, to avoid potential random IO. This closes #11687.
-
由 Zhijiang 提交于
-
由 Zhijiang 提交于
-
由 Zhijiang 提交于
-
由 Zhijiang 提交于
[hotfix][network] Fix useless backlog value in BufferAndAvailability returned by RemoteInputChannel#getNextBuffer
-
由 Leonard Xu 提交于
This closes #11755
-
由 Jiangjie (Becket) Qin 提交于
1. Add CoordinatedOperatorFactory interface. 2. Rename SourceReaderOpertor to SourceOperator, add implementation and connect it to OperatorEventGateway. 3. Rename SourceReaderStreamTask to SourceOperatorStreamTask 4. Fix some bugs in StreamTaskMailboxTestHarness.
-
由 godfreyhe 提交于
[FLINK-17252][table] Add Table#execute api and support SELECT statement in TableEnvironment#executeSql This closes #12049
-
由 vthinkxie 提交于
This closes #12085.
-
- 11 5月, 2020 10 次提交
-
-
由 Dawid Wysakowicz 提交于
-
由 wangyang0918 提交于
-
由 wangyang0918 提交于
[FLINK-17416][e2e][k8s] Use fixed v1.16.9 because fabric8 kubernetes-client could not work with higher version under jdk 8u252
-
由 Gary Yao 提交于
Reduce visiblity from public to private for methods - PipelinedRegionComputeUtilTest#assertSameRegion(Set<SchedulingExecutionVertex>...) - PipelinedRegionComputeUtilTest#assertDistinctRegions(Set<SchedulingExecutionVertex>...) This closes #11929.
-
由 Gary Yao 提交于
[FLINK-17369][tests] Rename RestartPipelinedRegionFailoverStrategyBuildingTest to PipelinedRegionComputeUtilTest
-
由 Gary Yao 提交于
[FLINK-17369][tests] In RestartPipelinedRegionFailoverStrategyBuildingTest invoke PipelinedRegionComputeUtil directly RestartPipelinedRegionFailoverStrategyBuildingTest means to test the logic in PipelinedRegionComputeUtil. Currently, however, the pipelined regions are retrieved indirectly from RestartPipelinedRegionFailoverStrategy. This commit changes that by using the PipelinedRegionComputeUtil directly from the test.
-
由 huangxingbo 提交于
[FLINK-17567][python][release] Create a dedicated Python directory in release directory to place Python-related source and binary packages This closes #12030.
-
由 Jingsong Lee 提交于
This closes #12063
-
由 acqua.csq 提交于
This closes #12061
-
由 厉颖 提交于
-