http://git-wip-us.apache.org/repos/asf/nifi/blob/96ed405d/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/iterator/SequentialRecordReaderEventIterator.java ---------------------------------------------------------------------- diff --git a/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/iterator/SequentialRecordReaderEventIterator.java b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/iterator/SequentialRecordReaderEventIterator.java new file mode 100644 index 0000000..869febf --- /dev/null +++ b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/iterator/SequentialRecordReaderEventIterator.java @@ -0,0 +1,115 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.nifi.provenance.store.iterator; + +import java.io.EOFException; +import java.io.File; +import java.io.FileNotFoundException; +import java.io.IOException; +import java.util.Collections; +import java.util.Iterator; +import java.util.List; +import java.util.Optional; + +import org.apache.nifi.provenance.ProvenanceEventRecord; +import org.apache.nifi.provenance.serialization.RecordReader; +import org.apache.nifi.provenance.store.RecordReaderFactory; + +public class SequentialRecordReaderEventIterator implements EventIterator { + private final Iterator<File> fileIterator; + private final RecordReaderFactory readerFactory; + private final long minimumEventId; + private final int maxAttributeChars; + + private boolean closed = false; + private RecordReader reader; + + public SequentialRecordReaderEventIterator(final List<File> filesToRead, final RecordReaderFactory readerFactory, final long minimumEventId, final int maxAttributeChars) { + this.fileIterator = filesToRead.iterator(); + this.readerFactory = readerFactory; + this.minimumEventId = minimumEventId; + this.maxAttributeChars = maxAttributeChars; + } + + @Override + public void close() throws IOException { + closed = true; + + if (reader != null) { + reader.close(); + } + } + + @Override + public Optional<ProvenanceEventRecord> nextEvent() throws IOException { + if (closed) { + throw new IOException("EventIterator is already closed"); + } + + if (reader == null) { + if (!rotateReader()) { + return Optional.empty(); + } + } + + while (true) { + final ProvenanceEventRecord event = reader.nextRecord(); + if (event == null) { + if (rotateReader()) { + continue; + } else { + return Optional.empty(); + } + } else { + return Optional.of(event); + } + } + } + + private boolean rotateReader() throws IOException { + final boolean readerExists = (reader != null); + if (readerExists) { + reader.close(); + } + + boolean multipleReadersOpened = false; + while (true) { + if (!fileIterator.hasNext()) { + return false; + } + + final File eventFile = fileIterator.next(); + try { + reader = readerFactory.newRecordReader(eventFile, Collections.emptyList(), maxAttributeChars); + break; + } catch (final FileNotFoundException | EOFException e) { + multipleReadersOpened = true; + // File may have aged off or was not fully written. Move to next file + continue; + } + } + + // If this is the first file in our list, the event of interest may not be the first event, + // so skip to the event that we want. + if (!readerExists && !multipleReadersOpened) { + reader.skipToEvent(minimumEventId); + } + + return true; + } +}
http://git-wip-us.apache.org/repos/asf/nifi/blob/96ed405d/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/toc/StandardTocReader.java ---------------------------------------------------------------------- diff --git a/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/toc/StandardTocReader.java b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/toc/StandardTocReader.java index 60328fa..a9c0f20 100644 --- a/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/toc/StandardTocReader.java +++ b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/toc/StandardTocReader.java @@ -16,12 +16,13 @@ */ package org.apache.nifi.provenance.toc; -import java.io.DataInputStream; import java.io.EOFException; import java.io.File; import java.io.FileInputStream; import java.io.IOException; +import org.apache.nifi.stream.io.StreamUtils; + /** * Standard implementation of TocReader. * @@ -38,27 +39,29 @@ public class StandardTocReader implements TocReader { private final boolean compressed; private final long[] offsets; private final long[] firstEventIds; + private final File file; public StandardTocReader(final File file) throws IOException { - try (final FileInputStream fis = new FileInputStream(file); - final DataInputStream dis = new DataInputStream(fis)) { + this.file = file; + final long fileLength = file.length(); + if (fileLength < 2) { + throw new EOFException(); + } - final int version = dis.read(); - if ( version < 0 ) { - throw new EOFException(); - } + try (final FileInputStream fis = new FileInputStream(file)) { + final byte[] buffer = new byte[(int) fileLength]; + StreamUtils.fillBuffer(fis, buffer); - final int compressionFlag = dis.read(); - if ( compressionFlag < 0 ) { - throw new EOFException(); - } + final int version = buffer[0]; + final int compressionFlag = buffer[1]; if ( compressionFlag == 0 ) { compressed = false; } else if ( compressionFlag == 1 ) { compressed = true; } else { - throw new IOException("Table of Contents appears to be corrupt: could not read 'compression flag' from header; expected value of 0 or 1 but got " + compressionFlag); + throw new IOException("Table of Contents file " + file + " appears to be corrupt: could not read 'compression flag' from header; " + + "expected value of 0 or 1 but got " + compressionFlag); } final int blockInfoBytes; @@ -72,7 +75,7 @@ public class StandardTocReader implements TocReader { break; } - final int numBlocks = (int) ((file.length() - 2) / blockInfoBytes); + final int numBlocks = (buffer.length - 2) / blockInfoBytes; offsets = new long[numBlocks]; if ( version > 1 ) { @@ -81,22 +84,41 @@ public class StandardTocReader implements TocReader { firstEventIds = new long[0]; } + int index = 2; for (int i=0; i < numBlocks; i++) { - offsets[i] = dis.readLong(); + offsets[i] = readLong(buffer, index); + index += 8; if ( version > 1 ) { - firstEventIds[i] = dis.readLong(); + firstEventIds[i] = readLong(buffer, index); + index += 8; } } } } + private long readLong(final byte[] buffer, final int offset) { + return ((long) buffer[offset] << 56) + + ((long) (buffer[offset + 1] & 0xFF) << 48) + + ((long) (buffer[offset + 2] & 0xFF) << 40) + + ((long) (buffer[offset + 3] & 0xFF) << 32) + + ((long) (buffer[offset + 4] & 0xFF) << 24) + + ((long) (buffer[offset + 5] & 0xFF) << 16) + + ((long) (buffer[offset + 6] & 0xFF) << 8) + + (buffer[offset + 7] & 0xFF); + } + @Override public boolean isCompressed() { return compressed; } @Override + public File getFile() { + return file; + } + + @Override public long getBlockOffset(final int blockIndex) { if ( blockIndex >= offsets.length ) { return -1L; @@ -105,6 +127,15 @@ public class StandardTocReader implements TocReader { } @Override + public long getFirstEventIdForBlock(final int blockIndex) { + if (blockIndex >= firstEventIds.length) { + return -1L; + } + + return firstEventIds[blockIndex]; + } + + @Override public long getLastBlockOffset() { if ( offsets.length == 0 ) { return 0L; @@ -113,7 +144,7 @@ public class StandardTocReader implements TocReader { } @Override - public void close() throws IOException { + public void close() { } @Override @@ -152,4 +183,9 @@ public class StandardTocReader implements TocReader { // Therefore, if the event is present, it must be in the last block. return firstEventIds.length - 1; } + + @Override + public String toString() { + return "StandardTocReader[file=" + file + ", compressed=" + compressed + "]"; + } } http://git-wip-us.apache.org/repos/asf/nifi/blob/96ed405d/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/toc/TocReader.java ---------------------------------------------------------------------- diff --git a/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/toc/TocReader.java b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/toc/TocReader.java index be6a165..0bb630e 100644 --- a/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/toc/TocReader.java +++ b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/toc/TocReader.java @@ -17,6 +17,7 @@ package org.apache.nifi.provenance.toc; import java.io.Closeable; +import java.io.File; /** * <p> @@ -37,15 +38,30 @@ public interface TocReader extends Closeable { boolean isCompressed(); /** + * @return the file that holds the TOC information + */ + File getFile(); + + /** * Returns the byte offset into the Journal File for the Block with the given index. * * @param blockIndex the block index to get the byte offset for * @return the byte offset for the given block index, or <code>-1</code> if the given block index - * does not exist + * does not exist */ long getBlockOffset(int blockIndex); /** + * Returns the ID of the first event that is found in the block with the given index, or -1 if + * the given block index does not exist. + * + * @param blockIndex the block index to get the first event id for + * @return the ID of the first event that is found in the block with the given index, or -1 if + * the given block index does not exist + */ + long getFirstEventIdForBlock(int blockIndex); + + /** * Returns the byte offset into the Journal File of the last Block in the given index * @return the byte offset into the Journal File of the last Block in the given index */ http://git-wip-us.apache.org/repos/asf/nifi/blob/96ed405d/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/toc/TocUtil.java ---------------------------------------------------------------------- diff --git a/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/toc/TocUtil.java b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/toc/TocUtil.java index 28cecd8..91f70db 100644 --- a/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/toc/TocUtil.java +++ b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/toc/TocUtil.java @@ -32,7 +32,7 @@ public class TocUtil { */ public static File getTocFile(final File journalFile) { final File tocDir = new File(journalFile.getParentFile(), "toc"); - final String basename = LuceneUtil.substringBefore(journalFile.getName(), "."); + final String basename = LuceneUtil.substringBefore(journalFile.getName(), ".prov"); final File tocFile = new File(tocDir, basename + ".toc"); return tocFile; } http://git-wip-us.apache.org/repos/asf/nifi/blob/96ed405d/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/util/CloseableUtil.java ---------------------------------------------------------------------- diff --git a/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/util/CloseableUtil.java b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/util/CloseableUtil.java new file mode 100644 index 0000000..26caa10 --- /dev/null +++ b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/util/CloseableUtil.java @@ -0,0 +1,45 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.nifi.provenance.util; + +import java.io.Closeable; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +public class CloseableUtil { + private static final Logger logger = LoggerFactory.getLogger(CloseableUtil.class); + + public static void closeQuietly(final Closeable... closeables) { + for (final Closeable closeable : closeables) { + if (closeable == null) { + continue; + } + + try { + closeable.close(); + } catch (final Exception e) { + if (logger.isDebugEnabled()) { + logger.warn("Failed to close {}; sources resources may not be cleaned up appropriately.", closeable, e); + } else { + logger.warn("Failed to close {}; sources resources may not be cleaned up appropriately.", closeable); + } + } + } + } +} http://git-wip-us.apache.org/repos/asf/nifi/blob/96ed405d/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/util/DirectoryUtils.java ---------------------------------------------------------------------- diff --git a/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/util/DirectoryUtils.java b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/util/DirectoryUtils.java new file mode 100644 index 0000000..a90500d --- /dev/null +++ b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/util/DirectoryUtils.java @@ -0,0 +1,96 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.nifi.provenance.util; + +import java.io.File; +import java.io.FileFilter; +import java.nio.file.Path; +import java.util.Arrays; +import java.util.Comparator; +import java.util.List; +import java.util.stream.Collectors; +import java.util.stream.Stream; + +import org.apache.nifi.provenance.RepositoryConfiguration; + +public class DirectoryUtils { + + public static final FileFilter EVENT_FILE_FILTER = f -> f.getName().endsWith(".prov") || f.getName().endsWith(".prov.gz"); + public static final FileFilter INDEX_FILE_FILTER = f -> f.getName().startsWith("index-"); + public static final Comparator<File> SMALLEST_ID_FIRST = (a, b) -> Long.compare(getMinId(a), getMinId(b)); + public static final Comparator<File> LARGEST_ID_FIRST = SMALLEST_ID_FIRST.reversed(); + public static final Comparator<File> OLDEST_INDEX_FIRST = (a, b) -> Long.compare(getIndexTimestamp(a), getIndexTimestamp(b)); + public static final Comparator<File> NEWEST_INDEX_FIRST = OLDEST_INDEX_FIRST.reversed(); + + public static List<Path> getProvenanceEventFiles(final RepositoryConfiguration repoConfig) { + return repoConfig.getStorageDirectories().values().stream() + .flatMap(f -> { + final File[] eventFiles = f.listFiles(EVENT_FILE_FILTER); + return eventFiles == null ? Stream.empty() : Arrays.stream(eventFiles); + }) + .map(f -> f.toPath()) + .collect(Collectors.toList()); + } + + public static long getMinId(final File file) { + final String filename = file.getName(); + final int firstDotIndex = filename.indexOf("."); + if (firstDotIndex < 1) { + return -1L; + } + + final String firstEventId = filename.substring(0, firstDotIndex); + try { + return Long.parseLong(firstEventId); + } catch (final NumberFormatException nfe) { + return -1L; + } + } + + public static long getIndexTimestamp(final File file) { + final String filename = file.getName(); + if (!filename.startsWith("index-") && filename.length() > 6) { + return -1L; + } + + final String suffix = filename.substring(6); + try { + return Long.parseLong(suffix); + } catch (final NumberFormatException nfe) { + return -1L; + } + } + + public static long getSize(final File file) { + if (file.isFile()) { + return file.length(); + } + + final File[] children = file.listFiles(); + if (children == null || children.length == 0) { + return 0L; + } + + long total = 0L; + for (final File child : children) { + total += getSize(child); + } + + return total; + } +} http://git-wip-us.apache.org/repos/asf/nifi/blob/96ed405d/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/util/DumpEventFile.java ---------------------------------------------------------------------- diff --git a/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/util/DumpEventFile.java b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/util/DumpEventFile.java new file mode 100644 index 0000000..df16356 --- /dev/null +++ b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/util/DumpEventFile.java @@ -0,0 +1,79 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.nifi.provenance.util; + +import java.io.File; +import java.io.IOException; +import java.util.Collections; +import java.util.Date; + +import org.apache.nifi.provenance.ProvenanceEventRecord; +import org.apache.nifi.provenance.StandardProvenanceEventRecord; +import org.apache.nifi.provenance.serialization.RecordReader; +import org.apache.nifi.provenance.serialization.RecordReaders; + +public class DumpEventFile { + + private static void printUsage() { + System.out.println("Usage:"); + System.out.println(); + System.out.println("java " + DumpEventFile.class.getName() + " <Event File to Dump>"); + System.out.println(); + } + + public static void main(final String[] args) throws IOException { + if (args.length != 1) { + printUsage(); + return; + } + + final File file = new File(args[0]); + if (!file.exists()) { + System.out.println("Cannot find file " + file.getAbsolutePath()); + return; + } + + try (final RecordReader reader = RecordReaders.newRecordReader(file, Collections.emptyList(), 65535)) { + StandardProvenanceEventRecord event; + int index = 0; + while ((event = reader.nextRecord()) != null) { + final long byteOffset = reader.getBytesConsumed(); + final String string = stringify(event, index++, byteOffset); + System.out.println(string); + } + } + } + + private static String stringify(final ProvenanceEventRecord event, final int index, final long byteOffset) { + final StringBuilder sb = new StringBuilder(); + sb.append("Event Index in File = ").append(index).append(", Byte Offset = ").append(byteOffset); + sb.append("\n\t").append("Event ID = ").append(event.getEventId()); + sb.append("\n\t").append("Event Type = ").append(event.getEventType()); + sb.append("\n\t").append("Event Time = ").append(new Date(event.getEventTime())); + sb.append("\n\t").append("Event UUID = ").append(event.getFlowFileUuid()); + sb.append("\n\t").append("Component ID = ").append(event.getComponentId()); + sb.append("\n\t").append("Event ID = ").append(event.getComponentType()); + sb.append("\n\t").append("Transit URI = ").append(event.getTransitUri()); + sb.append("\n\t").append("Parent IDs = ").append(event.getParentUuids()); + sb.append("\n\t").append("Child IDs = ").append(event.getChildUuids()); + sb.append("\n\t").append("Previous Attributes = ").append(event.getPreviousAttributes()); + sb.append("\n\t").append("Updated Attributes = ").append(event.getUpdatedAttributes()); + + return sb.toString(); + } +} http://git-wip-us.apache.org/repos/asf/nifi/blob/96ed405d/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/util/NamedThreadFactory.java ---------------------------------------------------------------------- diff --git a/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/util/NamedThreadFactory.java b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/util/NamedThreadFactory.java new file mode 100644 index 0000000..2ee6ed6 --- /dev/null +++ b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/util/NamedThreadFactory.java @@ -0,0 +1,47 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.nifi.provenance.util; + +import java.util.concurrent.Executors; +import java.util.concurrent.ThreadFactory; +import java.util.concurrent.atomic.AtomicInteger; + +public class NamedThreadFactory implements ThreadFactory { + + private final AtomicInteger counter = new AtomicInteger(0); + private final ThreadFactory defaultThreadFactory = Executors.defaultThreadFactory(); + private final String namePrefix; + private final boolean daemon; + + public NamedThreadFactory(final String namePrefix) { + this(namePrefix, false); + } + + public NamedThreadFactory(final String namePrefix, final boolean daemon) { + this.namePrefix = namePrefix; + this.daemon = daemon; + } + + @Override + public Thread newThread(final Runnable r) { + final Thread thread = defaultThreadFactory.newThread(r); + thread.setName(namePrefix + "-" + counter.incrementAndGet()); + thread.setDaemon(daemon); + return thread; + } +} http://git-wip-us.apache.org/repos/asf/nifi/blob/96ed405d/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/util/StorageSummaryEvent.java ---------------------------------------------------------------------- diff --git a/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/util/StorageSummaryEvent.java b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/util/StorageSummaryEvent.java new file mode 100644 index 0000000..41d5ade --- /dev/null +++ b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/util/StorageSummaryEvent.java @@ -0,0 +1,185 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.nifi.provenance.util; + +import java.util.List; +import java.util.Map; + +import org.apache.nifi.provenance.ProvenanceEventRecord; +import org.apache.nifi.provenance.ProvenanceEventType; +import org.apache.nifi.provenance.serialization.StorageSummary; + +public class StorageSummaryEvent implements ProvenanceEventRecord { + private final ProvenanceEventRecord event; + private final StorageSummary storageSummary; + + public StorageSummaryEvent(final ProvenanceEventRecord event, final StorageSummary storageSummary) { + this.event = event; + this.storageSummary = storageSummary; + } + + @Override + public long getEventId() { + return storageSummary.getEventId(); + } + + @Override + public long getEventTime() { + return event.getEventTime(); + } + + @Override + public long getFlowFileEntryDate() { + return event.getFlowFileEntryDate(); + } + + @Override + public long getLineageStartDate() { + return event.getLineageStartDate(); + } + + @Override + public long getFileSize() { + return event.getFileSize(); + } + + @Override + public Long getPreviousFileSize() { + return event.getPreviousFileSize(); + } + + @Override + public long getEventDuration() { + return event.getEventDuration(); + } + + @Override + public ProvenanceEventType getEventType() { + return event.getEventType(); + } + + @Override + public Map<String, String> getAttributes() { + return event.getAttributes(); + } + + @Override + public Map<String, String> getPreviousAttributes() { + return event.getPreviousAttributes(); + } + + @Override + public Map<String, String> getUpdatedAttributes() { + return event.getUpdatedAttributes(); + } + + @Override + public String getComponentId() { + return event.getComponentId(); + } + + @Override + public String getComponentType() { + return event.getComponentType(); + } + + @Override + public String getTransitUri() { + return event.getTransitUri(); + } + + @Override + public String getSourceSystemFlowFileIdentifier() { + return event.getSourceSystemFlowFileIdentifier(); + } + + @Override + public String getFlowFileUuid() { + return event.getFlowFileUuid(); + } + + @Override + public List<String> getParentUuids() { + return event.getParentUuids(); + } + + @Override + public List<String> getChildUuids() { + return event.getChildUuids(); + } + + @Override + public String getAlternateIdentifierUri() { + return event.getAlternateIdentifierUri(); + } + + @Override + public String getDetails() { + return event.getDetails(); + } + + @Override + public String getRelationship() { + return event.getRelationship(); + } + + @Override + public String getSourceQueueIdentifier() { + return event.getSourceQueueIdentifier(); + } + + @Override + public String getContentClaimSection() { + return event.getContentClaimSection(); + } + + @Override + public String getPreviousContentClaimSection() { + return event.getPreviousContentClaimSection(); + } + + @Override + public String getContentClaimContainer() { + return event.getContentClaimContainer(); + } + + @Override + public String getPreviousContentClaimContainer() { + return event.getPreviousContentClaimContainer(); + } + + @Override + public String getContentClaimIdentifier() { + return event.getContentClaimIdentifier(); + } + + @Override + public String getPreviousContentClaimIdentifier() { + return event.getPreviousContentClaimIdentifier(); + } + + @Override + public Long getContentClaimOffset() { + return event.getContentClaimOffset(); + } + + @Override + public Long getPreviousContentClaimOffset() { + return event.getPreviousContentClaimOffset(); + } +} http://git-wip-us.apache.org/repos/asf/nifi/blob/96ed405d/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/resources/META-INF/services/org.apache.nifi.provenance.ProvenanceRepository ---------------------------------------------------------------------- diff --git a/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/resources/META-INF/services/org.apache.nifi.provenance.ProvenanceRepository b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/resources/META-INF/services/org.apache.nifi.provenance.ProvenanceRepository index 78da70e..6a353d2 100644 --- a/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/resources/META-INF/services/org.apache.nifi.provenance.ProvenanceRepository +++ b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/resources/META-INF/services/org.apache.nifi.provenance.ProvenanceRepository @@ -12,4 +12,5 @@ # WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. # See the License for the specific language governing permissions and # limitations under the License. -org.apache.nifi.provenance.PersistentProvenanceRepository \ No newline at end of file +org.apache.nifi.provenance.PersistentProvenanceRepository +org.apache.nifi.provenance.WriteAheadProvenanceRepository \ No newline at end of file http://git-wip-us.apache.org/repos/asf/nifi/blob/96ed405d/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/test/java/org/apache/nifi/provenance/AbstractTestRecordReaderWriter.java ---------------------------------------------------------------------- diff --git a/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/test/java/org/apache/nifi/provenance/AbstractTestRecordReaderWriter.java b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/test/java/org/apache/nifi/provenance/AbstractTestRecordReaderWriter.java index bae2364..36397c4 100644 --- a/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/test/java/org/apache/nifi/provenance/AbstractTestRecordReaderWriter.java +++ b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/test/java/org/apache/nifi/provenance/AbstractTestRecordReaderWriter.java @@ -17,8 +17,8 @@ package org.apache.nifi.provenance; -import static org.apache.nifi.provenance.TestUtil.createFlowFile; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertNull; import static org.junit.Assert.assertTrue; @@ -27,8 +27,10 @@ import java.io.File; import java.io.FileInputStream; import java.io.IOException; import java.io.InputStream; -import java.util.HashMap; +import java.util.ArrayList; +import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.UUID; import org.apache.nifi.provenance.serialization.RecordReader; @@ -42,6 +44,7 @@ import org.apache.nifi.util.file.FileUtils; import org.junit.BeforeClass; import org.junit.Test; + public abstract class AbstractTestRecordReaderWriter { @BeforeClass public static void setLogLevel() { @@ -49,20 +52,7 @@ public abstract class AbstractTestRecordReaderWriter { } protected ProvenanceEventRecord createEvent() { - final Map<String, String> attributes = new HashMap<>(); - attributes.put("filename", "1.txt"); - attributes.put("uuid", UUID.randomUUID().toString()); - - final ProvenanceEventBuilder builder = new StandardProvenanceEventRecord.Builder(); - builder.setEventTime(System.currentTimeMillis()); - builder.setEventType(ProvenanceEventType.RECEIVE); - builder.setTransitUri("nifi://unit-test"); - builder.fromFlowFile(createFlowFile(3L, 3000L, attributes)); - builder.setComponentId("1234"); - builder.setComponentType("dummy processor"); - final ProvenanceEventRecord record = builder.build(); - - return record; + return TestUtil.createEvent(); } @Test @@ -73,7 +63,7 @@ public abstract class AbstractTestRecordReaderWriter { final RecordWriter writer = createWriter(journalFile, tocWriter, false, 1024 * 1024); writer.writeHeader(1L); - writer.writeRecord(createEvent(), 1L); + writer.writeRecord(createEvent()); writer.close(); final TocReader tocReader = new StandardTocReader(tocFile); @@ -101,7 +91,7 @@ public abstract class AbstractTestRecordReaderWriter { final RecordWriter writer = createWriter(journalFile, tocWriter, true, 8192); writer.writeHeader(1L); - writer.writeRecord(createEvent(), 1L); + writer.writeRecord(createEvent()); writer.close(); final TocReader tocReader = new StandardTocReader(tocFile); @@ -131,7 +121,7 @@ public abstract class AbstractTestRecordReaderWriter { writer.writeHeader(1L); for (int i = 0; i < 10; i++) { - writer.writeRecord(createEvent(), i); + writer.writeRecord(createEvent()); } writer.close(); @@ -170,7 +160,7 @@ public abstract class AbstractTestRecordReaderWriter { writer.writeHeader(1L); for (int i = 0; i < 10; i++) { - writer.writeRecord(createEvent(), i); + writer.writeRecord(createEvent()); } writer.close(); @@ -198,6 +188,56 @@ public abstract class AbstractTestRecordReaderWriter { FileUtils.deleteFile(journalFile.getParentFile(), true); } + @Test + public void testSkipToEvent() throws IOException { + final File journalFile = new File("target/storage/" + UUID.randomUUID().toString() + "/testSimpleWrite.gz"); + final File tocFile = TocUtil.getTocFile(journalFile); + final TocWriter tocWriter = new StandardTocWriter(tocFile, false, false); + // new block each 10 bytes + final RecordWriter writer = createWriter(journalFile, tocWriter, true, 100); + + writer.writeHeader(0L); + final int numEvents = 10; + final List<ProvenanceEventRecord> events = new ArrayList<>(); + for (int i = 0; i < numEvents; i++) { + final ProvenanceEventRecord event = createEvent(); + events.add(event); + writer.writeRecord(event); + } + writer.close(); + + final TocReader tocReader = new StandardTocReader(tocFile); + + try (final FileInputStream fis = new FileInputStream(journalFile); + final RecordReader reader = createReader(fis, journalFile.getName(), tocReader, 2048)) { + + for (int i = 0; i < numEvents; i++) { + final Optional<ProvenanceEventRecord> eventOption = reader.skipToEvent(i); + assertTrue(eventOption.isPresent()); + assertEquals(i, eventOption.get().getEventId()); + assertEquals(events.get(i), eventOption.get()); + + final StandardProvenanceEventRecord consumedEvent = reader.nextRecord(); + assertEquals(eventOption.get(), consumedEvent); + } + + assertFalse(reader.skipToEvent(numEvents + 1).isPresent()); + } + + try (final FileInputStream fis = new FileInputStream(journalFile); + final RecordReader reader = createReader(fis, journalFile.getName(), tocReader, 2048)) { + + for (int i = 0; i < 3; i++) { + final Optional<ProvenanceEventRecord> eventOption = reader.skipToEvent(8); + assertTrue(eventOption.isPresent()); + assertEquals(events.get(8), eventOption.get()); + } + + final StandardProvenanceEventRecord consumedEvent = reader.nextRecord(); + assertEquals(events.get(8), consumedEvent); + } + } + protected abstract RecordWriter createWriter(File file, TocWriter tocWriter, boolean compressed, int uncompressedBlockSize) throws IOException; protected abstract RecordReader createReader(InputStream in, String journalFilename, TocReader tocReader, int maxAttributeSize) throws IOException; http://git-wip-us.apache.org/repos/asf/nifi/blob/96ed405d/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/test/java/org/apache/nifi/provenance/TestEventIdFirstSchemaRecordReaderWriter.java ---------------------------------------------------------------------- diff --git a/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/test/java/org/apache/nifi/provenance/TestEventIdFirstSchemaRecordReaderWriter.java b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/test/java/org/apache/nifi/provenance/TestEventIdFirstSchemaRecordReaderWriter.java new file mode 100644 index 0000000..9c89ab3 --- /dev/null +++ b/nifi-nar-bundles/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/test/java/org/apache/nifi/provenance/TestEventIdFirstSchemaRecordReaderWriter.java @@ -0,0 +1,477 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.nifi.provenance; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertNull; + +import java.io.File; +import java.io.FileInputStream; +import java.io.IOException; +import java.io.InputStream; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.UUID; +import java.util.concurrent.Callable; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; + +import org.apache.nifi.provenance.serialization.RecordReader; +import org.apache.nifi.provenance.serialization.RecordWriter; +import org.apache.nifi.provenance.toc.StandardTocReader; +import org.apache.nifi.provenance.toc.StandardTocWriter; +import org.apache.nifi.provenance.toc.TocReader; +import org.apache.nifi.provenance.toc.TocUtil; +import org.apache.nifi.provenance.toc.TocWriter; +import org.apache.nifi.util.file.FileUtils; +import org.junit.Assert; +import org.junit.Before; +import org.junit.BeforeClass; +import org.junit.Ignore; +import org.junit.Test; + +public class TestEventIdFirstSchemaRecordReaderWriter extends AbstractTestRecordReaderWriter { + private final AtomicLong idGenerator = new AtomicLong(0L); + private File journalFile; + private File tocFile; + + @BeforeClass + public static void setupLogger() { + System.setProperty("org.slf4j.simpleLogger.log.org.apache.nifi", "DEBUG"); + } + + @Before + public void setup() { + journalFile = new File("target/storage/" + UUID.randomUUID().toString() + "/testEventIdFirstSchemaRecordReaderWriter"); + tocFile = TocUtil.getTocFile(journalFile); + idGenerator.set(0L); + }; + + @Test + public void testContentClaimUnchanged() throws IOException { + final File journalFile = new File("target/storage/" + UUID.randomUUID().toString() + "/testSimpleWrite.gz"); + final File tocFile = TocUtil.getTocFile(journalFile); + final TocWriter tocWriter = new StandardTocWriter(tocFile, false, false); + final RecordWriter writer = createWriter(journalFile, tocWriter, true, 8192); + + final Map<String, String> attributes = new HashMap<>(); + attributes.put("filename", "1.txt"); + attributes.put("uuid", UUID.randomUUID().toString()); + + final ProvenanceEventBuilder builder = new StandardProvenanceEventRecord.Builder(); + builder.setEventTime(System.currentTimeMillis()); + builder.setEventType(ProvenanceEventType.RECEIVE); + builder.setTransitUri("nifi://unit-test"); + builder.fromFlowFile(TestUtil.createFlowFile(3L, 3000L, attributes)); + builder.setComponentId("1234"); + builder.setComponentType("dummy processor"); + builder.setPreviousContentClaim("container-1", "section-1", "identifier-1", 1L, 1L); + builder.setCurrentContentClaim("container-1", "section-1", "identifier-1", 1L, 1L); + final ProvenanceEventRecord record = builder.build(); + + writer.writeHeader(1L); + writer.writeRecord(record); + writer.close(); + + final TocReader tocReader = new StandardTocReader(tocFile); + + try (final FileInputStream fis = new FileInputStream(journalFile); + final RecordReader reader = createReader(fis, journalFile.getName(), tocReader, 2048)) { + assertEquals(0, reader.getBlockIndex()); + reader.skipToBlock(0); + final StandardProvenanceEventRecord recovered = reader.nextRecord(); + assertNotNull(recovered); + + assertEquals("nifi://unit-test", recovered.getTransitUri()); + + assertEquals("container-1", recovered.getPreviousContentClaimContainer()); + assertEquals("container-1", recovered.getContentClaimContainer()); + + assertEquals("section-1", recovered.getPreviousContentClaimSection()); + assertEquals("section-1", recovered.getContentClaimSection()); + + assertEquals("identifier-1", recovered.getPreviousContentClaimIdentifier()); + assertEquals("identifier-1", recovered.getContentClaimIdentifier()); + + assertEquals(1L, recovered.getPreviousContentClaimOffset().longValue()); + assertEquals(1L, recovered.getContentClaimOffset().longValue()); + + assertEquals(1L, recovered.getPreviousFileSize().longValue()); + assertEquals(1L, recovered.getContentClaimOffset().longValue()); + + assertNull(reader.nextRecord()); + } + + FileUtils.deleteFile(journalFile.getParentFile(), true); + } + + @Test + public void testContentClaimRemoved() throws IOException { + final File journalFile = new File("target/storage/" + UUID.randomUUID().toString() + "/testSimpleWrite.gz"); + final File tocFile = TocUtil.getTocFile(journalFile); + final TocWriter tocWriter = new StandardTocWriter(tocFile, false, false); + final RecordWriter writer = createWriter(journalFile, tocWriter, true, 8192); + + final Map<String, String> attributes = new HashMap<>(); + attributes.put("filename", "1.txt"); + attributes.put("uuid", UUID.randomUUID().toString()); + + final ProvenanceEventBuilder builder = new StandardProvenanceEventRecord.Builder(); + builder.setEventTime(System.currentTimeMillis()); + builder.setEventType(ProvenanceEventType.RECEIVE); + builder.setTransitUri("nifi://unit-test"); + builder.fromFlowFile(TestUtil.createFlowFile(3L, 3000L, attributes)); + builder.setComponentId("1234"); + builder.setComponentType("dummy processor"); + builder.setPreviousContentClaim("container-1", "section-1", "identifier-1", 1L, 1L); + builder.setCurrentContentClaim(null, null, null, 0L, 0L); + final ProvenanceEventRecord record = builder.build(); + + writer.writeHeader(1L); + writer.writeRecord(record); + writer.close(); + + final TocReader tocReader = new StandardTocReader(tocFile); + + try (final FileInputStream fis = new FileInputStream(journalFile); + final RecordReader reader = createReader(fis, journalFile.getName(), tocReader, 2048)) { + assertEquals(0, reader.getBlockIndex()); + reader.skipToBlock(0); + final StandardProvenanceEventRecord recovered = reader.nextRecord(); + assertNotNull(recovered); + + assertEquals("nifi://unit-test", recovered.getTransitUri()); + + assertEquals("container-1", recovered.getPreviousContentClaimContainer()); + assertNull(recovered.getContentClaimContainer()); + + assertEquals("section-1", recovered.getPreviousContentClaimSection()); + assertNull(recovered.getContentClaimSection()); + + assertEquals("identifier-1", recovered.getPreviousContentClaimIdentifier()); + assertNull(recovered.getContentClaimIdentifier()); + + assertEquals(1L, recovered.getPreviousContentClaimOffset().longValue()); + assertNull(recovered.getContentClaimOffset()); + + assertEquals(1L, recovered.getPreviousFileSize().longValue()); + assertEquals(0L, recovered.getFileSize()); + + assertNull(reader.nextRecord()); + } + + FileUtils.deleteFile(journalFile.getParentFile(), true); + } + + @Test + public void testContentClaimAdded() throws IOException { + final File journalFile = new File("target/storage/" + UUID.randomUUID().toString() + "/testSimpleWrite.gz"); + final File tocFile = TocUtil.getTocFile(journalFile); + final TocWriter tocWriter = new StandardTocWriter(tocFile, false, false); + final RecordWriter writer = createWriter(journalFile, tocWriter, true, 8192); + + final Map<String, String> attributes = new HashMap<>(); + attributes.put("filename", "1.txt"); + attributes.put("uuid", UUID.randomUUID().toString()); + + final ProvenanceEventBuilder builder = new StandardProvenanceEventRecord.Builder(); + builder.setEventTime(System.currentTimeMillis()); + builder.setEventType(ProvenanceEventType.RECEIVE); + builder.setTransitUri("nifi://unit-test"); + builder.fromFlowFile(TestUtil.createFlowFile(3L, 3000L, attributes)); + builder.setComponentId("1234"); + builder.setComponentType("dummy processor"); + builder.setCurrentContentClaim("container-1", "section-1", "identifier-1", 1L, 1L); + final ProvenanceEventRecord record = builder.build(); + + writer.writeHeader(1L); + writer.writeRecord(record); + writer.close(); + + final TocReader tocReader = new StandardTocReader(tocFile); + + try (final FileInputStream fis = new FileInputStream(journalFile); + final RecordReader reader = createReader(fis, journalFile.getName(), tocReader, 2048)) { + assertEquals(0, reader.getBlockIndex()); + reader.skipToBlock(0); + final StandardProvenanceEventRecord recovered = reader.nextRecord(); + assertNotNull(recovered); + + assertEquals("nifi://unit-test", recovered.getTransitUri()); + + assertEquals("container-1", recovered.getContentClaimContainer()); + assertNull(recovered.getPreviousContentClaimContainer()); + + assertEquals("section-1", recovered.getContentClaimSection()); + assertNull(recovered.getPreviousContentClaimSection()); + + assertEquals("identifier-1", recovered.getContentClaimIdentifier()); + assertNull(recovered.getPreviousContentClaimIdentifier()); + + assertEquals(1L, recovered.getContentClaimOffset().longValue()); + assertNull(recovered.getPreviousContentClaimOffset()); + + assertEquals(1L, recovered.getFileSize()); + assertNull(recovered.getPreviousContentClaimOffset()); + + assertNull(reader.nextRecord()); + } + + FileUtils.deleteFile(journalFile.getParentFile(), true); + } + + @Test + public void testContentClaimChanged() throws IOException { + final File journalFile = new File("target/storage/" + UUID.randomUUID().toString() + "/testSimpleWrite.gz"); + final File tocFile = TocUtil.getTocFile(journalFile); + final TocWriter tocWriter = new StandardTocWriter(tocFile, false, false); + final RecordWriter writer = createWriter(journalFile, tocWriter, true, 8192); + + final Map<String, String> attributes = new HashMap<>(); + attributes.put("filename", "1.txt"); + attributes.put("uuid", UUID.randomUUID().toString()); + + final ProvenanceEventBuilder builder = new StandardProvenanceEventRecord.Builder(); + builder.setEventTime(System.currentTimeMillis()); + builder.setEventType(ProvenanceEventType.RECEIVE); + builder.setTransitUri("nifi://unit-test"); + builder.fromFlowFile(TestUtil.createFlowFile(3L, 3000L, attributes)); + builder.setComponentId("1234"); + builder.setComponentType("dummy processor"); + builder.setPreviousContentClaim("container-1", "section-1", "identifier-1", 1L, 1L); + builder.setCurrentContentClaim("container-2", "section-2", "identifier-2", 2L, 2L); + final ProvenanceEventRecord record = builder.build(); + + writer.writeHeader(1L); + writer.writeRecord(record); + writer.close(); + + final TocReader tocReader = new StandardTocReader(tocFile); + + try (final FileInputStream fis = new FileInputStream(journalFile); + final RecordReader reader = createReader(fis, journalFile.getName(), tocReader, 2048)) { + assertEquals(0, reader.getBlockIndex()); + reader.skipToBlock(0); + final StandardProvenanceEventRecord recovered = reader.nextRecord(); + assertNotNull(recovered); + + assertEquals("nifi://unit-test", recovered.getTransitUri()); + + assertEquals("container-1", recovered.getPreviousContentClaimContainer()); + assertEquals("container-2", recovered.getContentClaimContainer()); + + assertEquals("section-1", recovered.getPreviousContentClaimSection()); + assertEquals("section-2", recovered.getContentClaimSection()); + + assertEquals("identifier-1", recovered.getPreviousContentClaimIdentifier()); + assertEquals("identifier-2", recovered.getContentClaimIdentifier()); + + assertEquals(1L, recovered.getPreviousContentClaimOffset().longValue()); + assertEquals(2L, recovered.getContentClaimOffset().longValue()); + + assertEquals(1L, recovered.getPreviousFileSize().longValue()); + assertEquals(2L, recovered.getContentClaimOffset().longValue()); + + assertNull(reader.nextRecord()); + } + + FileUtils.deleteFile(journalFile.getParentFile(), true); + } + + @Test + public void testEventIdAndTimestampCorrect() throws IOException { + final File journalFile = new File("target/storage/" + UUID.randomUUID().toString() + "/testSimpleWrite.gz"); + final File tocFile = TocUtil.getTocFile(journalFile); + final TocWriter tocWriter = new StandardTocWriter(tocFile, false, false); + final RecordWriter writer = createWriter(journalFile, tocWriter, true, 8192); + + final Map<String, String> attributes = new HashMap<>(); + attributes.put("filename", "1.txt"); + attributes.put("uuid", UUID.randomUUID().toString()); + + final long timestamp = System.currentTimeMillis() - 10000L; + + final StandardProvenanceEventRecord.Builder builder = new StandardProvenanceEventRecord.Builder(); + builder.setEventId(1_000_000); + builder.setEventTime(timestamp); + builder.setEventType(ProvenanceEventType.RECEIVE); + builder.setTransitUri("nifi://unit-test"); + builder.fromFlowFile(TestUtil.createFlowFile(3L, 3000L, attributes)); + builder.setComponentId("1234"); + builder.setComponentType("dummy processor"); + builder.setPreviousContentClaim("container-1", "section-1", "identifier-1", 1L, 1L); + builder.setCurrentContentClaim("container-2", "section-2", "identifier-2", 2L, 2L); + final ProvenanceEventRecord record = builder.build(); + + writer.writeHeader(500_000L); + writer.writeRecord(record); + writer.close(); + + final TocReader tocReader = new StandardTocReader(tocFile); + + try (final FileInputStream fis = new FileInputStream(journalFile); + final RecordReader reader = createReader(fis, journalFile.getName(), tocReader, 2048)) { + + final ProvenanceEventRecord event = reader.nextRecord(); + assertNotNull(event); + assertEquals(1_000_000L, event.getEventId()); + assertEquals(timestamp, event.getEventTime()); + assertNull(reader.nextRecord()); + } + + FileUtils.deleteFile(journalFile.getParentFile(), true); + } + + + @Test + public void testComponentIdInlineAndLookup() throws IOException { + final File journalFile = new File("target/storage/" + UUID.randomUUID().toString() + "/testSimpleWrite.prov"); + final File tocFile = TocUtil.getTocFile(journalFile); + final TocWriter tocWriter = new StandardTocWriter(tocFile, false, false); + + final IdentifierLookup lookup = new IdentifierLookup() { + @Override + public List<String> getQueueIdentifiers() { + return Collections.emptyList(); + } + + @Override + public List<String> getComponentTypes() { + return Collections.singletonList("unit-test-component-1"); + } + + @Override + public List<String> getComponentIdentifiers() { + return Collections.singletonList("1234"); + } + }; + + final RecordWriter writer = new EventIdFirstSchemaRecordWriter(journalFile, idGenerator, tocWriter, false, 1024 * 32, lookup); + + final Map<String, String> attributes = new HashMap<>(); + attributes.put("filename", "1.txt"); + attributes.put("uuid", UUID.randomUUID().toString()); + + final StandardProvenanceEventRecord.Builder builder = new StandardProvenanceEventRecord.Builder(); + builder.setEventId(1_000_000); + builder.setEventTime(System.currentTimeMillis()); + builder.setEventType(ProvenanceEventType.RECEIVE); + builder.setTransitUri("nifi://unit-test"); + builder.fromFlowFile(TestUtil.createFlowFile(3L, 3000L, attributes)); + builder.setComponentId("1234"); + builder.setComponentType("unit-test-component-2"); + builder.setPreviousContentClaim("container-1", "section-1", "identifier-1", 1L, 1L); + builder.setCurrentContentClaim("container-2", "section-2", "identifier-2", 2L, 2L); + + writer.writeHeader(500_000L); + writer.writeRecord(builder.build()); + + builder.setEventId(1_000_001L); + builder.setComponentId("4444"); + builder.setComponentType("unit-test-component-1"); + writer.writeRecord(builder.build()); + + writer.close(); + + final TocReader tocReader = new StandardTocReader(tocFile); + + try (final FileInputStream fis = new FileInputStream(journalFile); + final RecordReader reader = createReader(fis, journalFile.getName(), tocReader, 2048)) { + + ProvenanceEventRecord event = reader.nextRecord(); + assertNotNull(event); + assertEquals(1_000_000L, event.getEventId()); + assertEquals("1234", event.getComponentId()); + assertEquals("unit-test-component-2", event.getComponentType()); + + event = reader.nextRecord(); + assertNotNull(event); + assertEquals(1_000_001L, event.getEventId()); + assertEquals("4444", event.getComponentId()); + assertEquals("unit-test-component-1", event.getComponentType()); + + assertNull(reader.nextRecord()); + } + + FileUtils.deleteFile(journalFile.getParentFile(), true); + } + + @Override + protected RecordWriter createWriter(final File file, final TocWriter tocWriter, final boolean compressed, final int uncompressedBlockSize) throws IOException { + return new EventIdFirstSchemaRecordWriter(file, idGenerator, tocWriter, compressed, uncompressedBlockSize, IdentifierLookup.EMPTY); + } + + @Override + protected RecordReader createReader(final InputStream in, final String journalFilename, final TocReader tocReader, final int maxAttributeSize) throws IOException { + return new EventIdFirstSchemaRecordReader(in, journalFilename, tocReader, maxAttributeSize); + } + + @Test + @Ignore + public void testPerformanceOfRandomAccessReads() throws Exception { + journalFile = new File("target/storage/" + UUID.randomUUID().toString() + "/testPerformanceOfRandomAccessReads.gz"); + tocFile = TocUtil.getTocFile(journalFile); + + final int blockSize = 1024 * 32; + try (final RecordWriter writer = createWriter(journalFile, new StandardTocWriter(tocFile, true, false), true, blockSize)) { + writer.writeHeader(0L); + + for (int i = 0; i < 100_000; i++) { + writer.writeRecord(createEvent()); + } + } + + final long[] eventIds = new long[] { + 4, 80, 1024, 1025, 1026, 1027, 1028, 1029, 1030, 40_000, 80_000, 99_000 + }; + + boolean loopForever = true; + while (loopForever) { + final long start = System.nanoTime(); + for (int i = 0; i < 1000; i++) { + try (final InputStream in = new FileInputStream(journalFile); + final RecordReader reader = createReader(in, journalFile.getName(), new StandardTocReader(tocFile), 32 * 1024)) { + + for (final long id : eventIds) { + time(() -> { + reader.skipToEvent(id); + return reader.nextRecord(); + }, id); + } + } + } + + final long ms = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - start); + System.out.println(ms + " ms total"); + } + } + + private void time(final Callable<StandardProvenanceEventRecord> task, final long id) throws Exception { + final long start = System.nanoTime(); + final StandardProvenanceEventRecord event = task.call(); + Assert.assertNotNull(event); + Assert.assertEquals(id, event.getEventId()); + // System.out.println(event); + final long nanos = System.nanoTime() - start; + final long millis = TimeUnit.NANOSECONDS.toMillis(nanos); + // System.out.println("Took " + millis + " ms to " + taskDescription); + } +}
