From a6ab711b5fc958b5bf62bb6540b320a5c111e69b Mon Sep 17 00:00:00 2001 From: Mark Payne Date: Tue, 3 Oct 2017 09:44:54 -0400 Subject: [PATCH 1/2] NIFI-4439: When a Provenance Event File is rolled over, we were failing to close the resource before attempting to compress it. Fixed that. --- .../nifi/provenance/store/WriteAheadStorePartition.java | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/WriteAheadStorePartition.java b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/WriteAheadStorePartition.java index fde76f5063d0..9184891b3917 100644 --- a/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/WriteAheadStorePartition.java +++ b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/WriteAheadStorePartition.java @@ -257,6 +257,10 @@ private synchronized boolean tryRollover(final RecordWriterLease lease) throws I final boolean updated = eventWriterLeaseRef.compareAndSet(lease, updatedLease); if (updated) { + if (lease != null) { + lease.close(); + } + updatedWriter.writeHeader(nextEventId); synchronized (minEventIdToPathMap) { From 42bcdca91a1e1f5cc666908904b770bc314c7329 Mon Sep 17 00:00:00 2001 From: Mark Payne Date: Mon, 16 Oct 2017 11:00:27 -0400 Subject: [PATCH 2/2] NIFI-4439: Addressed threading bug that can occur when rolling over provenance record writer --- .../store/WriteAheadStorePartition.java | 63 ++++++++++--------- 1 file changed, 34 insertions(+), 29 deletions(-) diff --git a/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/WriteAheadStorePartition.java b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/WriteAheadStorePartition.java index 9184891b3917..05ea17cf646c 100644 --- a/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/WriteAheadStorePartition.java +++ b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/WriteAheadStorePartition.java @@ -253,42 +253,47 @@ private synchronized boolean tryRollover(final RecordWriterLease lease) throws I final long nextEventId = idGenerator.get(); final File updatedEventFile = new File(partitionDirectory, nextEventId + ".prov"); final RecordWriter updatedWriter = recordWriterFactory.createWriter(updatedEventFile, idGenerator, false, true); - final RecordWriterLease updatedLease = new RecordWriterLease(updatedWriter, config.getMaxEventFileCapacity(), config.getMaxEventFileCount()); - final boolean updated = eventWriterLeaseRef.compareAndSet(lease, updatedLease); - - if (updated) { - if (lease != null) { - lease.close(); - } + + // Synchronize on the writer to ensure that no other thread is able to obtain the writer and start writing events to it until after it has + // been fully initialized (i.e., the header has been written, etc.) + synchronized (updatedWriter) { + final RecordWriterLease updatedLease = new RecordWriterLease(updatedWriter, config.getMaxEventFileCapacity(), config.getMaxEventFileCount()); + final boolean updated = eventWriterLeaseRef.compareAndSet(lease, updatedLease); + + if (updated) { + if (lease != null) { + lease.close(); + } - updatedWriter.writeHeader(nextEventId); + updatedWriter.writeHeader(nextEventId); - synchronized (minEventIdToPathMap) { - minEventIdToPathMap.put(nextEventId, updatedEventFile); - } + synchronized (minEventIdToPathMap) { + minEventIdToPathMap.put(nextEventId, updatedEventFile); + } - if (config.isCompressOnRollover() && lease != null && lease.getWriter() != null) { - boolean offered = false; - while (!offered && !closed) { - try { - offered = filesToCompress.offer(lease.getWriter().getFile(), 1, TimeUnit.SECONDS); - } catch (final InterruptedException ie) { - Thread.currentThread().interrupt(); - throw new IOException("Interrupted while waiting to enqueue " + lease.getWriter().getFile() + " for compression"); + if (config.isCompressOnRollover() && lease != null && lease.getWriter() != null) { + boolean offered = false; + while (!offered && !closed) { + try { + offered = filesToCompress.offer(lease.getWriter().getFile(), 1, TimeUnit.SECONDS); + } catch (final InterruptedException ie) { + Thread.currentThread().interrupt(); + throw new IOException("Interrupted while waiting to enqueue " + lease.getWriter().getFile() + " for compression"); + } } } - } - return true; - } else { - try { - updatedWriter.close(); - } catch (final Exception e) { - logger.warn("Failed to close Record Writer {}; some resources may not be cleaned up properly.", updatedWriter, e); - } + return true; + } else { + try { + updatedWriter.close(); + } catch (final Exception e) { + logger.warn("Failed to close Record Writer {}; some resources may not be cleaned up properly.", updatedWriter, e); + } - updatedEventFile.delete(); - return false; + updatedEventFile.delete(); + return false; + } } }