You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.24 04:35:14