From 7b571ca3e1d971bc485848fe0f119bb304e921e2 Mon Sep 17 00:00:00 2001 From: "nitin.kashyap" Date: Mon, 29 Mar 2021 17:18:43 +0530 Subject: [PATCH] added concurrent update protection for the finishedQuery::iterator --- .../dynamicfilter/DynamicFilterService.java | 34 +++++++++++-------- 1 file changed, 19 insertions(+), 15 deletions(-) diff --git a/presto-main/src/main/java/io/prestosql/dynamicfilter/DynamicFilterService.java b/presto-main/src/main/java/io/prestosql/dynamicfilter/DynamicFilterService.java index c61a534ff..fe498b5ab 100644 --- a/presto-main/src/main/java/io/prestosql/dynamicfilter/DynamicFilterService.java +++ b/presto-main/src/main/java/io/prestosql/dynamicfilter/DynamicFilterService.java @@ -226,24 +226,26 @@ public class DynamicFilterService List handledQuery = new ArrayList<>(); StateMap mergedStateCollection = (StateMap) stateStoreProvider.getStateStore().getOrCreateStateCollection(DynamicFilterUtils.MERGED_DYNAMIC_FILTERS, MAP); // Clear registered dynamic filter tasks - for (String queryId : finishedQuery) { - Map filters = dynamicFilters.get(queryId); - if (filters != null) { - for (Entry entry : filters.entrySet()) { - String filterId = entry.getKey(); - clearPartialResults(filterId, queryId); - if (entry.getValue().isMerged()) { - String filterKey = createKey(DynamicFilterUtils.FILTERPREFIX, filterId, queryId); - mergedStateCollection.remove(filterKey); + synchronized (finishedQuery) { + for (String queryId : finishedQuery) { + Map filters = dynamicFilters.get(queryId); + if (filters != null) { + for (Entry entry : filters.entrySet()) { + String filterId = entry.getKey(); + clearPartialResults(filterId, queryId); + if (entry.getValue().isMerged()) { + String filterKey = createKey(DynamicFilterUtils.FILTERPREFIX, filterId, queryId); + mergedStateCollection.remove(filterKey); + } } } - } - dynamicFilters.remove(queryId); + dynamicFilters.remove(queryId); - cachedDynamicFilters.remove(queryId); - handledQuery.add(queryId); + cachedDynamicFilters.remove(queryId); + handledQuery.add(queryId); + } + finishedQuery.removeAll(handledQuery); } - finishedQuery.removeAll(handledQuery); } private static BloomFilter mergeBloomFilters(Collection partialBloomFilters) @@ -355,7 +357,9 @@ public class DynamicFilterService */ public void clearDynamicFiltersForQuery(String queryId) { - finishedQuery.add(queryId); + synchronized (finishedQuery) { + finishedQuery.add(queryId); + } } /**