Skip to content
Merged
Show file tree
Hide file tree
Changes from 6 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions checkstyle/import-control.xml
Original file line number Diff line number Diff line change
Expand Up @@ -314,6 +314,7 @@

<subpackage name="raft">
<allow pkg="org.apache.kafka.raft" />
<allow pkg="org.apache.kafka.snapshot" />
<allow pkg="org.apache.kafka.clients" />
<allow pkg="org.apache.kafka.common.config" />
<allow pkg="org.apache.kafka.common.message" />
Expand All @@ -325,6 +326,11 @@
<allow pkg="com.fasterxml.jackson" />
</subpackage>

<subpackage name="snapshot">
<allow pkg="org.apache.kafka.raft" />
<allow pkg="org.apache.kafka.common.record" />
</subpackage>

<subpackage name="connect">
<allow pkg="org.apache.kafka.common" />
<allow pkg="org.apache.kafka.connect.data" />
Expand Down
17 changes: 17 additions & 0 deletions core/src/main/scala/kafka/raft/KafkaMetadataLog.scala
Original file line number Diff line number Diff line change
Expand Up @@ -16,16 +16,21 @@
*/
package kafka.raft

import java.nio.file.NoSuchFileException
import java.util.Optional

import kafka.log.{AppendOrigin, Log}
import kafka.server.{FetchHighWatermark, FetchLogEnd}
import kafka.snapshot.KafkaSnapshotWriter
import org.apache.kafka.common.record.{MemoryRecords, Records}
import org.apache.kafka.common.{KafkaException, TopicPartition}
import org.apache.kafka.raft
import org.apache.kafka.raft.{LogAppendInfo, LogFetchInfo, LogOffsetMetadata, Isolation, ReplicatedLog}
import org.apache.kafka.snapshot.SnapshotWriter
import org.apache.kafka.snapshot.SnapshotReader

import scala.compat.java8.OptionConverters._
import kafka.snapshot.KafkaSnapshotReader

class KafkaMetadataLog(
log: Log,
Expand Down Expand Up @@ -141,6 +146,18 @@ class KafkaMetadataLog(
topicPartition
}

override def createSnapshot(snapshotId: raft.OffsetAndEpoch): SnapshotWriter = {
KafkaSnapshotWriter(log.dir.toPath, snapshotId)
}

override def readSnapshot(snapshotId: raft.OffsetAndEpoch): Optional[SnapshotReader] = {
try {
Optional.of(KafkaSnapshotReader(log.dir.toPath, snapshotId))
} catch {
case e: NoSuchFileException => Optional.empty()
}
}

override def close(): Unit = {
log.close()
}
Expand Down
71 changes: 71 additions & 0 deletions core/src/main/scala/kafka/snapshot/KafkaSnapshotReader.scala
Original file line number Diff line number Diff line change
@@ -0,0 +1,71 @@
/*
* 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 kafka.snapshot

import java.nio.ByteBuffer
import java.nio.file.Path
import java.util.{Iterator => JIterator}
import org.apache.kafka.common.record.RecordBatch
import org.apache.kafka.common.record.FileRecords
import org.apache.kafka.raft.OffsetAndEpoch
import org.apache.kafka.snapshot.SnapshotReader

final class KafkaSnapshotReader private (fileRecords: FileRecords, snapshotId: OffsetAndEpoch) extends SnapshotReader {
Comment thread
jsancio marked this conversation as resolved.
Outdated
def snapshotId(): OffsetAndEpoch = {
snapshotId
}

def sizeInBytes(): Long = {
fileRecords.sizeInBytes()
}

def iterator(): JIterator[RecordBatch] = {
new JIterator[RecordBatch] {
private[this] val iterator = fileRecords.batchIterator()

override def hasNext(): Boolean = {
iterator.hasNext()
}

override def next(): RecordBatch = {
iterator.next()
}
}
}

def read(buffer: ByteBuffer, position: Long): Int = {
fileRecords.channel.read(buffer, position)
}

def close(): Unit = {
fileRecords.close()
}
}

object KafkaSnapshotReader {
def apply(logDir: Path, snapshotId: OffsetAndEpoch): KafkaSnapshotReader = {
val fileRecords = FileRecords.open(
snapshotPath(logDir, snapshotId).toFile,
false, // mutable
true, // fileAlreadyExists
0, // initFileSize
false // preallocate
)

new KafkaSnapshotReader(fileRecords, snapshotId)
}
}
95 changes: 95 additions & 0 deletions core/src/main/scala/kafka/snapshot/KafkaSnapshotWriter.scala
Original file line number Diff line number Diff line change
@@ -0,0 +1,95 @@
/*
* 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 kafka.snapshot

import java.nio.ByteBuffer
import java.nio.channels.FileChannel
import java.nio.file.Files
import java.nio.file.Path
import java.nio.file.StandardCopyOption
import java.nio.file.StandardOpenOption
import kafka.utils.Logging
import org.apache.kafka.common.record.MemoryRecords
import org.apache.kafka.common.utils.Utils
import org.apache.kafka.raft.OffsetAndEpoch
import org.apache.kafka.snapshot.SnapshotWriter

final class KafkaSnapshotWriter(
path: Path,
channel: FileChannel,
snapshotId: OffsetAndEpoch
) extends SnapshotWriter with Logging {
private[this] var frozen = false

override def snapshotId(): OffsetAndEpoch = {
snapshotId
}

override def sizeInBytes(): Long = {
channel.size()
}

override def append(records: MemoryRecords): Int = {
if (frozen) {
throw new IllegalStateException(s"Append not supported. Snapshot is already frozen: id = $snapshotId; path = $path")
}

records.writeFullyTo(channel)
}

override def append(buffer: ByteBuffer): Unit = {
if (frozen) {
throw new IllegalStateException(s"Append not supported. Snapshot is already frozen: id = $snapshotId; path = $path")
}

Utils.writeFully(channel, buffer)
}

override def isFrozen(): Boolean = {
frozen
}

override def freeze(): Unit = {
channel.close()
frozen = true

// Set readonly and ignore the result
if (!path.toFile.setReadOnly()) {
info(s"Unable to change permission to readonly for internal snapshot file '$path'")
}

val destination = moveRename(path, snapshotId)
Files.move(path, destination, StandardCopyOption.ATOMIC_MOVE)
}

override def close(): Unit = {
channel.close()
Files.deleteIfExists(path)
}
}

object KafkaSnapshotWriter {
def apply(logDir: Path, snapshotId: OffsetAndEpoch): KafkaSnapshotWriter = {
val path = createTempFile(logDir, snapshotId)

new KafkaSnapshotWriter(
path,
FileChannel.open(path, Utils.mkSet(StandardOpenOption.WRITE, StandardOpenOption.APPEND)),
snapshotId
)
}
}
59 changes: 59 additions & 0 deletions core/src/main/scala/kafka/snapshot/package.scala
Original file line number Diff line number Diff line change
@@ -0,0 +1,59 @@
/*
* 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 kafka

import java.nio.file.Files
import java.nio.file.Path
import java.text.NumberFormat
import org.apache.kafka.raft.OffsetAndEpoch

package object snapshot {
private[this] val SnapshotDir = "snapshots"
private[this] val Suffix = ".snapshot"
private[this] val PartialSuffix = s"$Suffix.part"

def snapshotDir(logDir: Path): Path = {
logDir.resolve(SnapshotDir)
}

def snapshotPath(logDir: Path, snapshotId: OffsetAndEpoch): Path = {
snapshotDir(logDir).resolve(filenameFromSnapshotId(snapshotId) + Suffix)
}

def filenameFromSnapshotId(snapshotId: OffsetAndEpoch): String = {
val formatter = NumberFormat.getInstance()
formatter.setMinimumIntegerDigits(20)
formatter.setGroupingUsed(false)

formatter.format(snapshotId.offset) + "-" + formatter.format(snapshotId.epoch)
}

def moveRename(source: Path, snapshotId: OffsetAndEpoch): Path = {
source.resolveSibling(filenameFromSnapshotId(snapshotId) + Suffix)
}

def createTempFile(logDir: Path, snapshotId: OffsetAndEpoch): Path = {
val dir = snapshotDir(logDir)

// Create the snapshot directory if it doesn't exists
Files.createDirectories(dir)

val prefix = s"${filenameFromSnapshotId(snapshotId)}-"

Files.createTempFile(dir, prefix, PartialSuffix)
}
}
Loading