如何在Spring Boot微服务中使用WebHDFS REST API操作Hadoop集群?
在Spring Boot中集成WebHDFS实现HDFS核心操作
下面是一套轻量方案,直接基于WebHDFS REST API实现创建文件夹、上传/读取/删除数据等操作,无需引入笨重的Hadoop客户端依赖:
1. 项目依赖配置
只需要Spring Boot基础Web依赖即可,用于发起HTTP请求和处理响应:
<dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-json</artifactId> </dependency> </dependencies>
2. 配置WebHDFS地址
在application.yml中添加集群WebHDFS的基础地址:
hadoop: webhdfs: url: http://你的Namenode主机地址:50070/webhdfs/v1
3. 封装WebHDFS操作工具类
创建工具类封装所有核心操作,用RestTemplate发起REST请求:
import org.springframework.beans.factory.annotation.Value; import org.springframework.http.HttpEntity; import org.springframework.http.HttpHeaders; import org.springframework.http.HttpMethod; import org.springframework.http.ResponseEntity; import org.springframework.stereotype.Component; import org.springframework.web.client.RestTemplate; import org.springframework.web.util.UriComponentsBuilder; import java.io.InputStream; import java.util.Map; @Component public class WebHdfsTemplate { @Value("${hadoop.webhdfs.url}") private String webHdfsBaseUrl; private final RestTemplate restTemplate; public WebHdfsTemplate(RestTemplate restTemplate) { this.restTemplate = restTemplate; } // 创建文件夹 public boolean createDirectory(String path, String permission) { String url = UriComponentsBuilder.fromHttpUrl(webHdfsBaseUrl + path) .queryParam("op", "MKDIRS") .queryParam("permission", permission) .toUriString(); ResponseEntity<Map> response = restTemplate.exchange(url, HttpMethod.PUT, null, Map.class); return (boolean) response.getBody().get("boolean"); } // 上传小文件(大文件需改用分块上传API) public boolean uploadFile(String path, InputStream fileStream, String permission) { // 第一步:获取上传重定向地址 String redirectUrl = UriComponentsBuilder.fromHttpUrl(webHdfsBaseUrl + path) .queryParam("op", "CREATE") .queryParam("permission", permission) .queryParam("overwrite", "true") .toUriString(); ResponseEntity<Void> redirectResponse = restTemplate.exchange(redirectUrl, HttpMethod.PUT, null, Void.class); String uploadUrl = redirectResponse.getHeaders().getLocation().toString(); // 第二步:上传文件内容 HttpHeaders headers = new HttpHeaders(); headers.set(HttpHeaders.CONTENT_TYPE, "application/octet-stream"); HttpEntity<InputStream> requestEntity = new HttpEntity<>(fileStream, headers); ResponseEntity<Map> uploadResponse = restTemplate.exchange(uploadUrl, HttpMethod.PUT, requestEntity, Map.class); return (boolean) uploadResponse.getBody().get("boolean"); } // 读取文件内容 public InputStream readFile(String path) { String url = UriComponentsBuilder.fromHttpUrl(webHdfsBaseUrl + path) .queryParam("op", "OPEN") .toUriString(); ResponseEntity<InputStream> response = restTemplate.exchange(url, HttpMethod.GET, null, InputStream.class); return response.getBody(); } // 删除文件/文件夹 public boolean delete(String path, boolean recursive) { String url = UriComponentsBuilder.fromHttpUrl(webHdfsBaseUrl + path) .queryParam("op", "DELETE") .queryParam("recursive", String.valueOf(recursive)) .toUriString(); ResponseEntity<Map> response = restTemplate.exchange(url, HttpMethod.DELETE, null, Map.class); return (boolean) response.getBody().get("boolean"); } }
同时需要配置RestTemplate的Bean:
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.web.client.RestTemplate; @Configuration public class RestTemplateConfig { @Bean public RestTemplate restTemplate() { return new RestTemplate(); } }
4. 业务服务层(可选)
如果需要封装业务逻辑,可创建服务类:
import org.springframework.stereotype.Service; import java.io.InputStream; @Service public class HdfsService { private final WebHdfsTemplate webHdfsTemplate; public HdfsService(WebHdfsTemplate webHdfsTemplate) { this.webHdfsTemplate = webHdfsTemplate; } public boolean createDir(String dirPath) { return webHdfsTemplate.createDirectory(dirPath, "755"); } public boolean uploadSmallFile(String filePath, InputStream fileStream) { return webHdfsTemplate.uploadFile(filePath, fileStream, "644"); } public InputStream getFileContent(String filePath) { return webHdfsTemplate.readFile(filePath); } public boolean deletePath(String path) { return webHdfsTemplate.delete(path, true); } }
5. 控制器示例(暴露REST接口)
编写控制器对外提供操作接口:
import org.springframework.http.HttpHeaders; import org.springframework.http.HttpStatus; import org.springframework.http.MediaType; import org.springframework.http.ResponseEntity; import org.springframework.web.bind.annotation.*; import org.springframework.web.multipart.MultipartFile; import java.io.IOException; import java.io.InputStream; @RestController @RequestMapping("/hdfs") public class HdfsController { private final HdfsService hdfsService; public HdfsController(HdfsService hdfsService) { this.hdfsService = hdfsService; } @PutMapping("/dir") public ResponseEntity<String> createDir(@RequestParam String path) { boolean success = hdfsService.createDir(path); return success ? ResponseEntity.ok("文件夹创建成功") : ResponseEntity.status(HttpStatus.INTERNAL_SERVER_ERROR).body("创建失败"); } @PostMapping("/file") public ResponseEntity<String> uploadFile(@RequestParam String filePath, @RequestParam MultipartFile file) throws IOException { boolean success = hdfsService.uploadSmallFile(filePath, file.getInputStream()); return success ? ResponseEntity.ok("文件上传成功") : ResponseEntity.status(HttpStatus.INTERNAL_SERVER_ERROR).body("上传失败"); } @GetMapping("/file") public ResponseEntity<InputStream> downloadFile(@RequestParam String filePath) { InputStream inputStream = hdfsService.getFileContent(filePath); HttpHeaders headers = new HttpHeaders(); headers.setContentType(MediaType.APPLICATION_OCTET_STREAM); headers.setContentDispositionFormData("attachment", filePath.substring(filePath.lastIndexOf("/") + 1)); return new ResponseEntity<>(inputStream, headers, HttpStatus.OK); } @DeleteMapping("/path") public ResponseEntity<String> deletePath(@RequestParam String path) { boolean success = hdfsService.deletePath(path); return success ? ResponseEntity.ok("删除成功") : ResponseEntity.status(HttpStatus.INTERNAL_SERVER_ERROR).body("删除失败"); } }
关键注意事项
- 确保HDFS集群已开启WebHDFS:在
hdfs-site.xml中配置dfs.webhdfs.enabled=true,并重启Namenode和Datanode。 - 端口验证:默认WebHDFS端口为50070,若集群修改过端口需同步调整配置。
- 权限控制:操作HDFS的用户默认是WebHDFS服务端运行用户,可通过请求参数
user.name=xxx指定操作用户。 - 大文件处理:上述示例仅适用于小文件,大文件需使用WebHDFS分块上传API(通过
offset参数控制分块位置)。
内容的提问来源于stack exchange,提问作者Kate
相关产品推荐
相关产品推荐

