如何基于IgnitePath获取InputStream?Apache Ignite相关技术疑问
嘿,刚接触Ignite和Hadoop的文件系统适配确实容易绕晕,我来给你拆解一下这两个问题~
一、如何确定read()方法的参数?
先逐个解释HadoopIgfsSecondaryFileSystemPositionedReadable.read(long pos, byte[] buf, int off, int len)的每个参数,再给你一个实用的代码示例:
pos:要读取的文件起始位置,单位是字节。如果从文件开头读就传0;如果是续读,就传上次读取结束后的位置。你可以先调用IgfsSecondaryFileSystem.fileLength(IgfsPath path)获取文件总长度,这样就能知道pos的合法范围是0到fileLength-1。buf:用来存储读取数据的字节数组,需要提前初始化(比如new byte[8192],选4k/8k这类常见的缓冲区大小都可以)。off:数据写入buf时的起始偏移量,一般传0就好(意思是从buf的第一个位置开始存数据);如果buf里已经有部分数据需要保留,再传对应的偏移值。len:想要读取的字节数,注意不能超过buf的剩余可用空间(buf.length - off),也不能超过文件剩余未读的字节数(fileLength - pos),否则会抛出异常或只读到剩余数据。
示例:读取整个文件
// 初始化你的Secondary FileSystem实例和文件路径 IgfsSecondaryFileSystem secondaryFs = ...; IgfsPath targetPath = new IgfsPath("/hadoop/storage/your-file.txt"); long totalFileLen = secondaryFs.fileLength(targetPath); HadoopIgfsSecondaryFileSystemPositionedReadable readable = (HadoopIgfsSecondaryFileSystemPositionedReadable)secondaryFs.open(targetPath); byte[] buffer = new byte[8192]; long currentPosition = 0; int bytesRead; try { while (currentPosition < totalFileLen) { // 计算本次能读取的最大字节数,避免超出文件范围 int readLength = (int)Math.min(buffer.length, totalFileLen - currentPosition); bytesRead = readable.read(currentPosition, buffer, 0, readLength); if (bytesRead == -1) break; // 到达文件末尾 // 这里替换成你的数据处理逻辑 processReadData(buffer, 0, bytesRead); // 更新当前读取位置 currentPosition += bytesRead; } } catch (IOException e) { // 处理IO异常 e.printStackTrace(); } finally { readable.close(); }
二、有没有其他方式获取标准的InputStream?
当然有,给你两种更省心的方案:
方案1:自己包装成标准InputStream
你可以写一个简单的包装类,把HadoopIgfsSecondaryFileSystemPositionedReadable转换成普通的InputStream,内部维护当前读取位置,自动处理定位逻辑:
public class PositionedReadableInputStream extends InputStream { private final HadoopIgfsSecondaryFileSystemPositionedReadable readable; private long currentPos; private final long totalFileLen; public PositionedReadableInputStream(HadoopIgfsSecondaryFileSystemPositionedReadable readable, long totalFileLen) { this.readable = readable; this.totalFileLen = totalFileLen; this.currentPos = 0; } @Override public int read() throws IOException { byte[] singleByteBuf = new byte[1]; int readCount = readable.read(currentPos, singleByteBuf, 0, 1); if (readCount == -1) return -1; currentPos++; return singleByteBuf[0] & 0xFF; } @Override public int read(byte[] b, int off, int len) throws IOException { if (currentPos >= totalFileLen) return -1; // 计算实际可读取的字节数 int actualReadLen = (int)Math.min(len, totalFileLen - currentPos); int readCount = readable.read(currentPos, b, off, actualReadLen); if (readCount > 0) currentPos += readCount; return readCount; } @Override public void close() throws IOException { readable.close(); } }
使用时直接把定位可读对象传入包装类,就能像普通输入流一样操作:
// 省略初始化代码... long fileLen = secondaryFs.fileLength(targetPath); HadoopIgfsSecondaryFileSystemPositionedReadable readable = (HadoopIgfsSecondaryFileSystemPositionedReadable)secondaryFs.open(targetPath); try (InputStream standardIn = new PositionedReadableInputStream(readable, fileLen)) { byte[] buf = new byte[4096]; int bytesRead; while ((bytesRead = standardIn.read(buf)) != -1) { // 处理数据 } } catch (IOException e) { e.printStackTrace(); }
方案2:通过IgniteFileSystem获取标准输入流
如果你是在Ignite集群环境中操作,直接使用Ignite提供的IgniteFileSystem主接口,它的open(IgfsPath)方法会返回IgfsInputStream——这个类直接继承自java.io.InputStream,完全不用自己处理定位逻辑!
// 启动Ignite并获取FileSystem实例 Ignite ignite = Ignition.start("/path/to/your/ignite-config.xml"); IgniteFileSystem igfs = ignite.fileSystem("your-igfs-name"); // 对应配置中的IGFS名称 IgfsPath targetPath = new IgfsPath("/secondary-fs/your-file.txt"); try (InputStream in = igfs.open(targetPath)) { // 像操作普通文件流一样读取数据 byte[] buf = new byte[8192]; int bytesRead; while ((bytesRead = in.read(buf)) != -1) { processReadData(buf, 0, bytesRead); } } catch (IOException e) { e.printStackTrace(); } finally { ignite.close(); }
这个方案是最推荐的,因为Ignite已经帮你封装了二级存储的适配逻辑,直接用标准流API就好。
内容的提问来源于stack exchange,提问作者blizzfire
相关产品推荐
相关产品推荐

