added concurrent update protection for the finishedQuery::iterator

This commit is contained in:
nitin.kashyap 2021-03-29 17:18:43 +05:30
parent c03b77b861
commit 7b571ca3e1
1 changed files with 19 additions and 15 deletions

View File

@ -226,24 +226,26 @@ public class DynamicFilterService
List<String> handledQuery = new ArrayList<>();
StateMap mergedStateCollection = (StateMap) stateStoreProvider.getStateStore().getOrCreateStateCollection(DynamicFilterUtils.MERGED_DYNAMIC_FILTERS, MAP);
// Clear registered dynamic filter tasks
for (String queryId : finishedQuery) {
Map<String, DynamicFilterRegistryInfo> filters = dynamicFilters.get(queryId);
if (filters != null) {
for (Entry<String, DynamicFilterRegistryInfo> 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<String, DynamicFilterRegistryInfo> filters = dynamicFilters.get(queryId);
if (filters != null) {
for (Entry<String, DynamicFilterRegistryInfo> 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<Object> partialBloomFilters)
@ -355,7 +357,9 @@ public class DynamicFilterService
*/
public void clearDynamicFiltersForQuery(String queryId)
{
finishedQuery.add(queryId);
synchronized (finishedQuery) {
finishedQuery.add(queryId);
}
}
/**