DorisIOManager.java
// 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.doris.paimon;
import org.apache.paimon.disk.BufferFileReader;
import org.apache.paimon.disk.BufferFileWriter;
import org.apache.paimon.disk.FileIOChannel;
import org.apache.paimon.disk.IOManager;
import org.apache.paimon.memory.Buffer;
import java.io.File;
import java.io.IOException;
import java.io.UncheckedIOException;
import java.nio.channels.FileChannel;
import java.nio.file.Files;
import java.nio.file.NoSuchFileException;
import java.nio.file.Path;
import java.util.HashSet;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.stream.Stream;
/** Paimon IOManager adapter which charges temporary I/O to Doris spill management. */
final class DorisIOManager implements IOManager {
interface SpillAccountant {
String[] getSpillDirectories() throws IOException;
void reserve(String path, long bytes) throws IOException;
void rollback(String path, long bytes);
void commitWrite(String path, long bytes);
void recordRead(String path, long bytes);
void release(String path, long bytes);
}
static final class SpillDirectoryCleanupException extends IOException {
private SpillDirectoryCleanupException(Exception cause) {
super("Failed to eagerly clean a Paimon spill directory", cause);
}
}
private final SpillAccountant accountant;
private final Map<String, Long> channelBytes = new ConcurrentHashMap<>();
private final Map<String, Long> directFileBytes = new ConcurrentHashMap<>();
private final Set<String> directFilePaths = ConcurrentHashMap.newKeySet();
private final Set<String> directDirectoryPaths = ConcurrentHashMap.newKeySet();
private final Set<String> bufferFilePaths = ConcurrentHashMap.newKeySet();
private final Map<String, Integer> activeChannelWriters = new ConcurrentHashMap<>();
private volatile IOManager delegate;
static DorisIOManager create(long nativeSpillSession) {
return new DorisIOManager(new NativeSpillAccountant(nativeSpillSession));
}
DorisIOManager(SpillAccountant accountant) {
this(null, accountant);
}
DorisIOManager(IOManager delegate, SpillAccountant accountant) {
this.accountant = accountant;
this.delegate = delegate;
}
private IOManager delegate() throws IOException {
if (delegate == null) {
synchronized (this) {
if (delegate == null) {
String[] spillDirectories = accountant.getSpillDirectories();
if (spillDirectories == null || spillDirectories.length == 0) {
throw new IOException("Doris spill manager returned no available directories");
}
delegate = IOManager.create(spillDirectories);
}
}
}
return delegate;
}
private IOManager uncheckedDelegate() {
try {
return delegate();
} catch (IOException e) {
throw new UncheckedIOException("Failed to initialize the Doris spill directory", e);
}
}
@Override
public FileIOChannel.ID createChannel() {
return uncheckedDelegate().createChannel();
}
@Override
public FileIOChannel.ID createChannel(String prefix) {
FileIOChannel.ID channel = uncheckedDelegate().createChannel(prefix);
directFilePaths.add(absolutePath(channel.getPath()));
return channel;
}
@Override
public String[] tempDirs() {
return uncheckedDelegate().tempDirs();
}
@Override
public String pickTempDir() {
String tempDir = uncheckedDelegate().pickTempDir();
directDirectoryPaths.add(absolutePath(tempDir));
return tempDir;
}
@Override
public FileIOChannel.Enumerator createChannelEnumerator() {
return uncheckedDelegate().createChannelEnumerator();
}
@Override
public BufferFileWriter createBufferFileWriter(FileIOChannel.ID channelID) throws IOException {
// ExternalBuffer clears old buffer channels with File.delete(), bypassing IOManager
// deletion. Release those known channels before reserving more space.
releaseDeletedChannels();
String path = absolutePath(channelID.getPath());
bufferFilePaths.add(path);
directFilePaths.remove(path);
releaseDirectFile(path);
return new AccountingBufferFileWriter(delegate().createBufferFileWriter(channelID), this);
}
@Override
public BufferFileReader createBufferFileReader(FileIOChannel.ID channelID) throws IOException {
return new AccountingBufferFileReader(delegate().createBufferFileReader(channelID), this);
}
@Override
public void close() throws Exception {
IOManager initializedDelegate = delegate;
if (initializedDelegate == null) {
return;
}
IOException accountingFailure = null;
try {
reconcileDirectFileGrowth();
} catch (IOException e) {
accountingFailure = e;
}
try {
initializedDelegate.close();
} catch (Exception e) {
releaseDeletedChannels();
releaseDeletedDirectFiles();
if (accountingFailure != null) {
e.addSuppressed(accountingFailure);
}
throw new SpillDirectoryCleanupException(e);
}
releaseDeletedChannels();
releaseDeletedDirectFiles();
if (accountingFailure != null) {
throw accountingFailure;
}
}
/** Reconciles files Paimon writes directly through paths returned by this IOManager. */
synchronized void reconcileDirectFileGrowth() throws IOException {
IOManager initializedDelegate = delegate;
if (initializedDelegate == null) {
return;
}
Set<String> seen = new HashSet<>();
for (String directFilePath : directFilePaths) {
Path path = Path.of(directFilePath);
if (Files.isRegularFile(path)) {
reconcileDirectPath(path, seen);
}
}
for (String tempDir : directDirectoryPaths) {
Path root = Path.of(tempDir);
if (!Files.exists(root)) {
continue;
}
try (Stream<Path> paths = Files.walk(root)) {
for (Path path : (Iterable<Path>) paths.filter(Files::isRegularFile)::iterator) {
String absolutePath = absolutePath(path);
if (bufferFilePaths.contains(absolutePath)) {
continue;
}
reconcileDirectPath(path, seen);
}
}
}
for (Map.Entry<String, Long> entry : directFileBytes.entrySet()) {
if (!seen.contains(entry.getKey())
&& directFileBytes.remove(entry.getKey(), entry.getValue())) {
accountant.release(entry.getKey(), entry.getValue());
}
}
}
private void reconcileDirectPath(Path path, Set<String> seen) throws IOException {
String absolutePath = absolutePath(path);
try {
long bytes = Files.size(path);
seen.add(absolutePath);
reconcileDirectFile(absolutePath, bytes);
} catch (NoSuchFileException ignored) {
// Paimon compaction can delete an obsolete SST while the directory is scanned.
}
}
private void reconcileDirectFile(String path, long bytes) throws IOException {
long accounted = directFileBytes.getOrDefault(path, 0L);
long delta = bytes - accounted;
if (delta > 0) {
accountant.reserve(path, delta);
directFileBytes.put(path, bytes);
accountant.commitWrite(path, delta);
} else if (delta < 0) {
directFileBytes.put(path, bytes);
accountant.release(path, -delta);
}
}
private void releaseDirectFile(String path) {
Long released = directFileBytes.remove(path);
if (released != null) {
accountant.release(path, released);
}
}
private void releaseDeletedDirectFiles() {
for (Map.Entry<String, Long> entry : directFileBytes.entrySet()) {
if (!new File(entry.getKey()).exists()
&& directFileBytes.remove(entry.getKey(), entry.getValue())) {
accountant.release(entry.getKey(), entry.getValue());
}
}
}
private static String absolutePath(String path) {
return absolutePath(Path.of(path));
}
private static String absolutePath(Path path) {
return path.toAbsolutePath().normalize().toString();
}
private boolean releaseDeletedChannels() {
boolean releasedAny = false;
for (Map.Entry<String, Long> entry : channelBytes.entrySet()) {
if (!activeChannelWriters.containsKey(entry.getKey())
&& !new File(entry.getKey()).exists()
&& channelBytes.remove(entry.getKey(), entry.getValue())) {
accountant.release(entry.getKey(), entry.getValue());
releasedAny = true;
}
}
return releasedAny;
}
private void reserveWrite(FileIOChannel.ID channelID, long bytes) throws IOException {
String path = channelID.getPath();
activeChannelWriters.merge(path, 1, Integer::sum);
boolean accounted = false;
boolean tracked = false;
try {
try {
accountant.reserve(path, bytes);
} catch (IOException reserveFailure) {
if (!releaseDeletedChannels()) {
throw reserveFailure;
}
accountant.reserve(path, bytes);
}
accounted = true;
channelBytes.merge(path, bytes, Long::sum);
tracked = true;
} finally {
if (!tracked) {
if (accounted) {
accountant.rollback(path, bytes);
}
finishWrite(channelID);
}
}
}
private void finishWrite(FileIOChannel.ID channelID) {
activeChannelWriters.computeIfPresent(channelID.getPath(),
(ignored, writers) -> writers == 1 ? null : writers - 1);
}
private void releaseChannel(FileIOChannel.ID channelID) {
Long released = channelBytes.remove(channelID.getPath());
if (released != null) {
accountant.release(channelID.getPath(), released);
}
}
private static final class NativeSpillAccountant implements SpillAccountant {
private final long nativeSpillSession;
private NativeSpillAccountant(long nativeSpillSession) {
this.nativeSpillSession = nativeSpillSession;
}
@Override
public String[] getSpillDirectories() throws IOException {
return PaimonJniWriter.getPaimonSpillDirectories(nativeSpillSession);
}
@Override
public void reserve(String path, long bytes) throws IOException {
PaimonJniWriter.reservePaimonSpill(nativeSpillSession, path, bytes);
}
@Override
public void rollback(String path, long bytes) {
PaimonJniWriter.updatePaimonSpillAccounting(
nativeSpillSession, path, -bytes, 0, 0);
}
@Override
public void commitWrite(String path, long bytes) {
PaimonJniWriter.updatePaimonSpillAccounting(
nativeSpillSession, path, 0, bytes, 0);
}
@Override
public void recordRead(String path, long bytes) {
PaimonJniWriter.updatePaimonSpillAccounting(
nativeSpillSession, path, 0, 0, bytes);
}
@Override
public void release(String path, long bytes) {
PaimonJniWriter.updatePaimonSpillAccounting(
nativeSpillSession, path, -bytes, 0, 0);
}
}
private static final class AccountingBufferFileWriter implements BufferFileWriter {
private final BufferFileWriter delegate;
private final DorisIOManager manager;
private AccountingBufferFileWriter(BufferFileWriter delegate, DorisIOManager manager) {
this.delegate = delegate;
this.manager = manager;
}
@Override
public void writeBlock(Buffer buffer) throws IOException {
long bytes = Integer.BYTES + buffer.getSize();
manager.reserveWrite(getChannelID(), bytes);
try {
delegate.writeBlock(buffer);
manager.accountant.commitWrite(getChannelID().getPath(), bytes);
} catch (IOException | RuntimeException writeFailure) {
try {
delegate.closeAndDelete();
} catch (IOException | RuntimeException cleanupFailure) {
writeFailure.addSuppressed(cleanupFailure);
}
if (!getChannelID().getPathFile().exists()) {
manager.releaseChannel(getChannelID());
}
throw writeFailure;
} finally {
manager.finishWrite(getChannelID());
}
}
@Override
public FileIOChannel.ID getChannelID() {
return delegate.getChannelID();
}
@Override
public long getSize() throws IOException {
return delegate.getSize();
}
@Override
public boolean isClosed() {
return delegate.isClosed();
}
@Override
public void close() throws IOException {
delegate.close();
}
@Override
public void deleteChannel() {
try {
delegate.deleteChannel();
} finally {
if (!getChannelID().getPathFile().exists()) {
manager.releaseChannel(getChannelID());
}
}
}
@Override
public FileChannel getNioFileChannel() {
return delegate.getNioFileChannel();
}
@Override
public void closeAndDelete() throws IOException {
try {
delegate.closeAndDelete();
} finally {
if (!getChannelID().getPathFile().exists()) {
manager.releaseChannel(getChannelID());
}
}
}
}
private static final class AccountingBufferFileReader implements BufferFileReader {
private final BufferFileReader delegate;
private final DorisIOManager manager;
private AccountingBufferFileReader(BufferFileReader delegate, DorisIOManager manager) {
this.delegate = delegate;
this.manager = manager;
}
@Override
public void readInto(Buffer buffer) throws IOException {
long position = delegate.getNioFileChannel().position();
delegate.readInto(buffer);
manager.accountant.recordRead(
getChannelID().getPath(), delegate.getNioFileChannel().position() - position);
}
@Override
public boolean hasReachedEndOfFile() {
return delegate.hasReachedEndOfFile();
}
@Override
public FileIOChannel.ID getChannelID() {
return delegate.getChannelID();
}
@Override
public long getSize() throws IOException {
return delegate.getSize();
}
@Override
public boolean isClosed() {
return delegate.isClosed();
}
@Override
public void close() throws IOException {
delegate.close();
}
@Override
public void deleteChannel() {
try {
delegate.deleteChannel();
} finally {
if (!getChannelID().getPathFile().exists()) {
manager.releaseChannel(getChannelID());
}
}
}
@Override
public FileChannel getNioFileChannel() {
return delegate.getNioFileChannel();
}
@Override
public void closeAndDelete() throws IOException {
try {
delegate.closeAndDelete();
} finally {
if (!getChannelID().getPathFile().exists()) {
manager.releaseChannel(getChannelID());
}
}
}
}
}