Overview
The library needs enhanced features for configuration management, comprehensive error handling, logging, and performance optimization to make it production-ready.
Current State
- No configuration management system
- Basic error handling only
- No logging framework
- No performance monitoring
- Missing validation utilities
Required Implementation
1. Configuration Management
object DFSConfig {
case class Config(
defaultReplication: Short = 3,
defaultBlockSize: Long = 128 * 1024 * 1024, // 128MB
bufferSize: Int = 4096,
maxRetries: Int = 3,
retryDelayMs: Long = 1000,
timeoutMs: Long = 30000,
compressionCodec: String = "org.apache.hadoop.io.compress.SnappyCodec"
)
object Config {
def load(configPath: String): Config = {
val conf = new Configuration()
conf.addResource(new Path(configPath))
Config(
defaultReplication = conf.getShort("dfs.replication", 3),
defaultBlockSize = conf.getLongBytes("dfs.blocksize", 128 * 1024 * 1024),
bufferSize = conf.getInt("io.file.buffer.size", 4096),
maxRetries = conf.getInt("dfs.client.retry.max.attempts", 3),
retryDelayMs = conf.getLong("dfs.client.retry.interval.ms", 1000),
timeoutMs = conf.getLong("dfs.client.socket.timeout", 30000),
compressionCodec = conf.get("dfs.client.compression.codec", "org.apache.hadoop.io.compress.SnappyCodec")
)
}
def apply(): Config = load("dfs-site.xml")
}
}
2. Comprehensive Error Handling
package dfs.errors
sealed trait DFSError extends Exception {
def message: String
def cause: Option[Throwable]
override def getMessage: String = message
}
case class PathNotFoundError(path: String, cause: Option[Throwable] = None)
extends DFSError {
def message: String = s"Path not found: $path"
}
case class PermissionDeniedError(path: String, user: String, cause: Option[Throwable] = None)
extends DFSError {
def message: String = s"Permission denied for user '$user' on path: $path"
}
case class InvalidPathError(path: String, reason: String, cause: Option[Throwable] = None)
extends DFSError {
def message: String = s"Invalid path '$path': $reason"
}
case class FileAlreadyExistsError(path: String, cause: Option[Throwable] = None)
extends DFSError {
def message: String = s"File already exists: $path"
}
case class DirectoryNotEmptyError(path: String, cause: Option[Throwable] = None)
extends DFSError {
def message: String = s"Directory not empty: $path"
}
case class OperationTimeoutError(operation: String, timeoutMs: Long, cause: Option[Throwable] = None)
extends DFSError {
def message: String = s"Operation '$operation' timed out after ${timeoutMs}ms"
}
object ErrorHandler {
def withRetry[T](maxRetries: Int = 3, delayMs: Long = 1000)(operation: => T): T = {
var attempt = 0
var lastException: Throwable = null
while (attempt < maxRetries) {
try {
return operation
} catch {
case e: Exception if attempt < maxRetries - 1 =>
lastException = e
attempt += 1
Thread.sleep(delayMs)
case e: Exception =>
throw e
}
}
throw lastException
}
def handleFileSystemError[T](path: String)(operation: => T): Either[DFSError, T] = {
try {
Right(operation)
} catch {
case e: FileNotFoundException => Left(PathNotFoundError(path, Some(e)))
case e: AccessControlException => Left(PermissionDeniedError(path, System.getProperty("user.name"), Some(e)))
case e: IllegalArgumentException => Left(InvalidPathError(path, e.getMessage, Some(e)))
case e: FileAlreadyExistsException => Left(FileAlreadyExistsError(path, Some(e)))
case e: IOException if e.getMessage.contains("not empty") => Left(DirectoryNotEmptyError(path, Some(e)))
case e: Exception => Left(new DFSError {
def message: String = s"Unexpected error: ${e.getMessage}"
def cause: Option[Throwable] = Some(e)
})
}
}
}
3. Logging Framework
import org.slf4j.{Logger, LoggerFactory}
trait Logging {
protected val logger: Logger = LoggerFactory.getLogger(getClass.getName)
protected def logOperation[T](operation: String, path: String)(f: => T): T = {
logger.info(s"Starting $operation on path: $path")
val startTime = System.currentTimeMillis()
try {
val result = f
val duration = System.currentTimeMillis() - startTime
logger.info(s"Completed $operation on path: $path in ${duration}ms")
result
} catch {
case e: Exception =>
logger.error(s"Failed $operation on path: $path", e)
throw e
}
}
}
object DFSLogger {
def apply[T](operation: String, path: String)(f: => T): T = {
val logger = LoggerFactory.getLogger("DFSOperations")
logger.info(s"Starting $operation on path: $path")
val startTime = System.currentTimeMillis()
try {
val result = f
val duration = System.currentTimeMillis() - startTime
logger.info(s"Completed $operation on path: $path in ${duration}ms")
result
} catch {
case e: Exception =>
logger.error(s"Failed $operation on path: $path", e)
throw e
}
}
}
4. Validation Utilities
object PathValidator {
def validatePath(path: String): Either[InvalidPathError, String] = {
if (path == null || path.trim.isEmpty) {
Left(InvalidPathError(path, "Path cannot be null or empty"))
} else if (!path.startsWith("/")) {
Left(InvalidPathError(path, "Path must be absolute"))
} else if (path.contains("//")) {
Left(InvalidPathError(path, "Path contains double slashes"))
} else {
Right(path)
}
}
def validateDirectoryPath(fs: FileSystem, path: String): Either[DFSError, String] = {
validatePath(path).flatMap { validPath =>
val pathObj = new Path(validPath)
if (!fs.exists(pathObj)) {
Left(PathNotFoundError(validPath))
} else if (!fs.getFileStatus(pathObj).isDirectory) {
Left(InvalidPathError(validPath, "Path is not a directory"))
} else {
Right(validPath)
}
}
}
def validateFilePath(fs: FileSystem, path: String): Either[DFSError, String] = {
validatePath(path).flatMap { validPath =>
val pathObj = new Path(validPath)
if (!fs.exists(pathObj)) {
Left(PathNotFoundError(validPath))
} else if (fs.getFileStatus(pathObj).isDirectory) {
Left(InvalidPathError(validPath, "Path is a directory, expected file"))
} else {
Right(validPath)
}
}
}
}
5. Performance Monitoring
object PerformanceMonitor {
case class OperationMetrics(
operation: String,
path: String,
durationMs: Long,
bytesProcessed: Long,
timestamp: Long = System.currentTimeMillis()
)
private val metricsBuffer = new java.util.concurrent.ConcurrentLinkedQueue[OperationMetrics]()
def record[T](operation: String, path: String)(f: => T): T = {
val startTime = System.currentTimeMillis()
var bytesProcessed = 0L
try {
val result = f
bytesProcessed = calculateBytesProcessed(result)
result
} finally {
val duration = System.currentTimeMillis() - startTime
metricsBuffer.offer(OperationMetrics(operation, path, duration, bytesProcessed))
}
}
private def calculateBytesProcessed(result: Any): Long = result match {
case bytes: Array[Byte] => bytes.length
case status: FileStatus => status.getLen
case statuses: Array[FileStatus] => statuses.map(_.getLen).sum
case _ => 0L
}
def getMetrics: List[OperationMetrics] = {
import scala.jdk.CollectionConverters._
metricsBuffer.asScala.toList
}
def getAverageDuration(operation: String): Option[Double] = {
val metrics = getMetrics.filter(_.operation == operation)
if (metrics.nonEmpty) {
Some(metrics.map(_.durationMs.toDouble).sum / metrics.size)
} else None
}
}
Integration Strategy
- Create new package structure:
dfs.config, dfs.errors, dfs.logging, dfs.validation, dfs.monitoring
- Integrate error handling into all existing operations
- Add configuration loading to initialization
- Add logging to all operations using the Logging trait
- Add validation to all public APIs
Acceptance Criteria
Priority
High - These are essential for production readiness
Labels
enhancement, production-readiness, high-priority, configuration, error-handling, monitoring
Overview
The library needs enhanced features for configuration management, comprehensive error handling, logging, and performance optimization to make it production-ready.
Current State
Required Implementation
1. Configuration Management
2. Comprehensive Error Handling
3. Logging Framework
4. Validation Utilities
5. Performance Monitoring
Integration Strategy
dfs.config,dfs.errors,dfs.logging,dfs.validation,dfs.monitoringAcceptance Criteria
Priority
High - These are essential for production readiness
Labels
enhancement, production-readiness, high-priority, configuration, error-handling, monitoring