修复 Apache Pekko Connectors S3 与华为云存储 OBS 兼容性问题
背景#
在软件开发领域文件访问是绕不开的内容,无论是从最早的 FTP,还是到后来的 SFTP、WebDAV,乃至现在占据主流的 S3 兼容协议,文件存取一直在向着访问安全、接入简易和弹性扩展的方向演进。因此越来越多的云厂商或开源项目锚定 S3 协议不断开发各种衍生产品。
正好之前项目中涉及文件存取相关的需求,在适配和开发过程中遇到一些问题记录在此处。目前问题已经提交回官方并经过了一轮 Review,目前是待合并状态。
第一次封装#
首先,我们定义一个 IStorageService 接口,并为该接口提供多个实现类,在实际使用时则通过配置项动态决定使用何种实现。基于演示的目的实际代码做了简化。
public interface IStorageService {
Vendor getVendor();
InputStream downloadFile(String path);
}
考虑到可能接入的兼容 S3 实现我们选取了 AWS S3,MinIO 和 Huawei OBS,接下来以标准的 AWS S3 为例。
@Slf4j
@Service
@ConditionalOnProperty(prefix = "file", value = "vendor", havingValue = "aws", matchIfMissing = false)
public class AWSStorageServiceImpl implements IStorageService {
private final FileProperties properties;
private final S3Client s3Client;
public AWSStorageServiceImpl(FileProperties properties) {
this.properties = properties;
this.s3Client = S3Client.builder()
.region(Region.of(properties.getRegion()))
.credentialsProvider(() -> AwsBasicCredentials.create(properties.getAccessKey(), properties.getSecretKey()))
.build();
}
@Override
public Vendor getVendor() {
return Vendor.AWS;
}
@Override
public InputStream downloadFile(String path) {
try {
GetObjectRequest getObjectRequest = GetObjectRequest.builder()
.bucket(defaultBucket())
.key(path)
.build();
ResponseInputStream<GetObjectResponse> response = s3Client.getObject(getObjectRequest);
return response;
} catch (Exception e) {
log.error("Download file error - AWS - path: {} Exception: ", path, e);
}
return null;
}
}
这里我们只实现了一个 downloadFile,在实际使用时至少还需要提供 文件上传,获取元数据,创建桶 和 查看桶内文件列表 等功能。至少目前下载文件的功能看起来工作的不错,同时 单接口多实现 也充分考虑到了可能会对接多种 S3 兼容实现的场景。
使用 Pekko Connectors S3 进行统一#
第一版虽然针对不同存储厂商给出了不同实现可供选择,但作为实际项目需要同时维护多份实现基本相同的代码,使用时也需要运维人员根据实际情况进行配置,一定程度上提升了研发和运维成本。基于以上弊端我们使用 apache pekko-connectors-s3 进行统一改造。
@Slf4j
@Service
public class UNIStorageServiceImpl implements IStorageService {
private final ActorSystem actorSystem;
private final FileProperties properties;
private final ObjectMapper objectMapper;
public UNIStorageServiceImpl(FileProperties properties, ActorSystem actorSystem, ObjectMapper objectMapper) {
this.properties = properties;
this.actorSystem = actorSystem;
this.objectMapper = objectMapper;
}
@Override
public InputStream downloadFile(String path) {
try {
final Source<ByteString, CompletionStage<ObjectMetadata>> source = S3.getObject(defaultBucket(), path);
InputStream is = source.toMat(StreamConverters.asInputStream(Duration.ofSeconds(properties.getTimeout())), Keep.right())
.run(actorSystem);
return is;
} catch (Exception e) {
log.error("Download file error - Pekko - path: {} Exception: ", path, e);
}
return null;
}
}
在切换到 pekko-connectors-s3 后代码整体变得更加简洁,同时不再需要维护多份代码代码的可维护性更好了。
问题跟踪#
在实际使用时发现下载 Huawei OBS 中文件会出现 403 Forbidden 错误但在使用 MinIO 进行测试时并未遇到该问题:
java.util.concurrent.ExecutionException: org.apache.pekko.stream.connectors.s3.S3Exception: (Status code: 403 Forbidden, Code: 403 Forbidden, RequestId: -, Resource: -)
at java.base/java.util.concurrent.CompletableFuture.wrapInExecutionException(CompletableFuture.java:345)
at java.base/java.util.concurrent.CompletableFuture.reportGet(CompletableFuture.java:440)
at java.base/java.util.concurrent.CompletableFuture.get(CompletableFuture.java:2117)
at scala.concurrent.impl.FutureConvertersImpl$CF.super$get(FutureConvertersImpl.scala:89)
at scala.concurrent.impl.FutureConvertersImpl$CF.$anonfun$get$2(FutureConvertersImpl.scala:89)
at scala.concurrent.BlockContext$DefaultBlockContext$.blockOn(BlockContext.scala:62)
at scala.concurrent.impl.FutureConvertersImpl$CF.get(FutureConvertersImpl.scala:89)
从日志中看不到更多有效信息,还好 pekko-connectors-s3 提供了 forward-proxy 配置以便我们通过代理观察请求和响应情况:
# An address of a proxy that will be used for all connections using HTTP CONNECT tunnel.
# forward-proxy {
# scheme = "https"
# host = "proxy"
# port = 8080
# credentials {
# username = "username"
# password = "password"
# }
#
以下是下载文件时返回的错误:
<Error>
<Code>SignatureDoesNotMatch</Code>
<Message>The request signature we calculated does not match the signature you provided. Check your key and signing method.</Message>
<RequestId>000001A0ECBCF74297E74C2F476D53FC</RequestId>
<HostId>Z9v+cC1sRnaWw6x0vi8pxxYA0YVnKxbYHUPAFpnxkX8sLV44u5b02Z+ailn2wCnR</HostId>
<AWSAccessKeyId>HPUAJIJO7SLODDRWWGEH</AWSAccessKeyId>
<SignatureProvided>2484abdd29693104aaa1f6f565cdb73aaee0fb9fb94b24e7ff04149e4ca311b4</SignatureProvided>
<StringToSign>AWS4-HMAC-SHA256
20260929T103641Z 20260929/us-east-1/s3/aws4_request
4e4f0057bff748b235d8ad1f2081ea85f51e86f928e638fb7369e63c0a602c48</StringToSign>
<StringToSignBytes>41 57 53 34 2d 48 4d 41 43 2d 53 48 41 32 35 36 0a 32 30 32 36 30 39 32 39 54 31 30 33 36 34 31 5a 0a 32 30 32 36 30 39 32 39 2f 75 73 2d 65 61 73 74 2d 31 2f 73 33 2f 61 77 73 34 5f 72 65 71 75 65 73 74 0a 34 65 34 66 30 30 35 37 62 66 66 37 34 38 62 32 33 35 64 38 61 64 31 66 32 30 38 31 65 61 38 35 66 35 31 65 38 36 66 39 32 38 65 36 33 38 66 62 37 33 36 39 65 36 33 63 30 61 36 30 32 63 34 38</StringToSignBytes>
<CanonicalRequest>GET
/2026/1233456&amp;&amp;7&amp;&amp;8&amp;&amp;9.xml
host:obs.cn-east-3.myhuaweicloud.com x-amz-content-sha256:e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855
x-amz-date:20260929T103641Z host;x-amz-content-sha256;x-amz-date
e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855</CanonicalRequest>
<StringToSignBytes>47 45 54 0a 2f 32 30 32 36 2f 31 32 33 33 34 35 36 26 26 37 26 26 38 26 26 39 2e 78 6d 6c 0a 0a 68 6f 73 74 3a 6f 62 73 2e 63 6e 2d 65 61 73 74 2d 33 2e 6d 79 68 75 61 77 65 69 63 6c 6f 75 64 2e 63 6f 6d 0a 78 2d 61 6d 7a 2d 63 6f 6e 74 65 6e 74 2d 73 68 61 32 35 36 3a 65 33 62 30 63 34 34 32 39 38 66 63 31 63 31 34 39 61 66 62 66 34 63 38 39 39 36 66 62 39 32 34 32 37 61 65 34 31 65 34 36 34 39 62 39 33 34 63 61 34 39 35 39 39 31 62 37 38 35 32 62 38 35 35 0a 78 2d 61 6d 7a 2d 64 61 74 65 3a 32 30 32 36 30 39 32 39 54 31 30 33 36 34 31 5a 0a 0a 68 6f 73 74 3b 78 2d 61 6d 7a 2d 63 6f 6e 74 65 6e 74 2d 73 68 61 32 35 36 3b 78 2d 61 6d 7a 2d 64 61 74 65 0a 65 33 62 30 63 34 34 32 39 38 66 63 31 63 31 34 39 61 66 62 66 34 63 38 39 39 36 66 62 39 32 34 32 37 61 65 34 31 65 34 36 34 39 62 39 33 34 63 61 34 39 35 39 39 31 62 37 38 35 32 62 38 35 35</StringToSignBytes>
</Error>
从以上错误响应可以看出:问题在于请求签名校验失败了。首先,根据以上请求来看,第一个怀疑的是 AWSAccessKeyId 配置是否存在问题,但考虑到 Huawei OBS 既然声明兼容 AWS S3,那么这里应当不是问题;同样 RequestId 和 HostId 作为请求分析和排查标识通常不会作为签名校验的一部分,因此同样略过;接下来是 CanonicalRequest 和 CanonicalRequestBytes,这两部分明显是一个整体,由于其和具体文件绑定,我们稍后再看;先检查 StringToSign 和 StringToSignBytes 部分,这里可以看出 StringToSign 同样是一个互关联的配置并且其中包含了 Region 使用的是 us-east-1 明显与实际使用不一致。在修正了 Region 配置后问题并没有解决,接下来我们继续验证 403 问题是普遍性的问题还是与特定文件相关联的。在进一步处理之前可以先用以下 Python 代码将 Byte 字符串转为 String 以便直观对比:
byte_str = """some-bytes-here
"""
hex_parts = byte_str.split()
result = bytes.fromhex("".join(hex_parts)).decode("UTF-8")
print(result)
解析后的内容如下:
GET
/2026/1233456&&7&&8&&9.xml
host:obs.cn-east-3.myhuaweicloud.com
x-amz-content-sha256:e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855
x-amz-date:20260929T103641Z
host;x-amz-content-sha256;x-amz-date
e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855
可以看出与原始 CanonicalRequest 一致(除了因 XML 语法原因进行编码的 & 符号),此时基于对编解码的敏感性,应当可以大致确定问题是由编解码差异引起的。
经过验证发现 403 Forbidden 问题确实在文件中包含 & 时会出现,而对应常规仅含 数字 和 英文字母 以及 中文 等的文件名时一切正常。既然问题明确了,最简单的方案是从业务上对文件名进行规范,但考虑到使用原始 "com.huaweicloud" % "esdk-obs-java-bundle" % "3.26.6" 依赖时并不存在问题,另外还存在大量历史文件需要处理,因此有必要进一步明确关于请求签名生成的逻辑。
首先克隆 apache/pekko-connectors 最新代码并根据当前项目所使用的版本切换到对应分支或 tag,这里我们以 1.3.0 为例。依据之前针对失败响应分析出的信息可快速定位到代码 s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/auth/CanonicalRequest.scala,浏览该文件可以发现有以下代码片段:
// https://tools.ietf.org/html/rfc3986#section-2.2
// Excludes "/" as it is an exception according to spec.
val reservedCharacters: String = ":?#[]@!$&'()*+,;="
显然,这里定义的保留字符中便包含此前 403 问题中文件名所使用的字符串,为进一步确认是否是保留字符的问题,我们手动修改 reservedCharacters 定义,并移除 & 字符后重新打包。
经验证,此时已经可以正常下载了。
验证#
虽然 Huawei OBS 的问题得以解决,但我们需要反思:为何使用 MinIO 作为测试环境时没有遇到该问题?因此,在 MinIO 和 Huawei OBS 之外,又选取了 Garage、Cloudflare R2、腾讯云 - 对象存储 COS 和 七牛云 - 对象存储 Kodo 等 S3 实现进行测试。
| vendor | 部署方式 | 保留字符串行为 |
|---|---|---|
| MinIO | 私有 | 一致 |
| Garage | 私有 | 一致 |
| CloudFlare R2 | 云厂商 | 一致 |
| Huawei OBS | 云厂商 | 不一致 |
| Tencent COS | 云厂商 | 不一致 |
| Qiniu Kodo | 云厂商 | 一致 |
从验证结果看,对于保留字符串的处理行为与 RFC 规范不一致的都是国内的存储厂商,而国外和开源的 MinIO 和 Garage 都遵守标准 RFC 行为。
合并回上游#
然问题已经解决,但是我们应该更多地将精力放在业务演进上,而不是维护一个专有分支。当然,最简单的方式就是在 Github 上提 issue 即可,这里我们多走一步直接把改好的代码通过 fork 和 pull request 形式提交到上游仓库。
将代码合并回上游时硬编码的方式就不再合适了。首先我们看一下 apache/pekko-connectors 仓库中关于贡献和代码准则是如何约定的。只有充分遵守社区规范才能更快地将我们的代码合并回主仓库。根据项目中 README.md 中的描述可以看到该仓库关于贡献者和代码行为准则都在 pekko-connectors/blob/main/CONTRIBUTING.md 和 Code of Conduct 两个文件中。其中我们需要重点关注的是关于编码风格、单元测试、commit 注释风格等相关内容,在熟悉完以上内容后就可以开始着手修改代码了。首先第一步是 fork 主仓库,接下来就是将所有修改提交到自己 fork 下来的仓库和分支中,完成 commit 提交后再到对应主仓库中发起 pull request,并描述清楚为什么要将修改合并回主仓库即可。
Round 1#
此前我们通过硬编码方式明确并解决了通过 华为云 对象存储服务 (OBS) 出现 Status code: 403 Forbidden 问题。
考虑到 Apache Pekko Connectors S3 作为通用的 S3 绑定,在对接不同 S3 兼容存储时,不一定所有实现针对保留字符串的行为都一致,因此,我们将保留字部分以可配置项的形式暴露出来,以便使用者可以根据实际情况自行选择。
首先,我们在 s3/src/main/resources/reference.conf 文件中添加关于 reserved-characters 的配置:
pekko.connectors.s3 {
...
reserved-characters = ":?#[]@!$&'()*+,;="
}
里默认使用 RFC 3986 中规定的保留字符。后续调用者在使用时如果有特殊需求,需要修改保留字符,只需在 application.conf 文件中自定义即可。
接下来是修改 s3/src/main/scala/org/apache/pekko/stream/connectors/s3/settings.scala 文件以便实现配置的读取:
val reservedCharacters: Option[String]
这里将 reservedCharacters 定义为 Option 类型是因为一般情况下不需要额外处理保留字,仅在接入那些未完全遵守 RFC 3986 规范厂的商服务时才需要显式配置。
/** Java API */
def getReservedCharacters: java.util.Optional[String] = reservedCharacters.toJava
val reservedCharacters = if (c.hasPath("reserved-characters")) {
Option(c.getString("reserved-characters"))
} else {
None
}
之后是修改 s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/auth/CanonicalRequest.scala 中关于保留字符的判断逻辑:
def isReservedCharacter(c: Char)(implicit conf: S3Settings): Boolean = {
if (conf.reservedCharacters.isDefined)
conf.getReservedCharacters.get().contains(c)
else
reservedCharacters.contains(c)
}
这里使用了 implicit隐式变量避免显式传参问题,类似于 Spring 的自动装配。
最后不能忘了单元测试:
def getSettings(
bufferType: BufferType = MemoryBufferType,
awsCredentials: AwsCredentialsProvider = AnonymousCredentialsProvider.create(),
s3Region: Region = Region.US_EAST_1,
listBucketApiVersion: ApiVersion = ApiVersion.ListBucketVersion2) = {
val regionProvider = new AwsRegionProvider {
def getRegion = s3Region
}
S3Settings(bufferType, awsCredentials, regionProvider, listBucketApiVersion, Map.empty)
}
implicit val settings: S3Settings = getSettings()
由于我们之前重新实现了 CanonicalRequest 中 def isReservedCharacter 方法,同时在方法定义上追加了一个隐式变量,因此这里单元测试只需要增加一个 S3Settings 的隐式变量定义即可。在 Scala 中通常使用 Scala Test,其行为更偏向 BDD 风格:
it should "only encode the configured reserved characters, leaving other default ones unencoded" in {
// '(' and ')' are in the default set, but not in the custom one
val customSettings = getSettings(reservedCharacters = Some("!"))
val req = pathRequest(Uri.Path.Empty / "file-(1)!.txt")
val canonicalRequest = CanonicalRequest.from(req)(customSettings)
canonicalRequest.canonicalString should equal(expectedCanonical("/file-(1)%21.txt"))
(CanonicalRequest.from(req).canonicalString should not).equal(canonicalRequest.canonicalString)
}
以上只列出一小段的单元测试,除了 CanonicalRequest 其他所有有修改逻辑地方最好都补充单元测试。
Round 2#
上一版基本上完成了功能拓展并考虑到了多厂商适配,但代码更多地还是以 Java 风格书写,并不太符合 Scala 风格,因此接下来我们使用更地道的 Scala 风格改进以上代码。
首先是 isReservedCharacter 方法,之前在第一版实现中我们采用了 if-else 条件分支,但在 Scala 中如果使用 Option 的话则更推荐使用模式匹配,因此我们将实现进一步改成以下形式:
def isReservedCharacter(c: Char)(implicit conf: S3Settings): Boolean =
conf.reservedCharacters match {
case Some(chars) => chars.contains(c)
case None => reservedCharacters.contains(c)
}
另外就是在 settings.scala 中我们定义的 reservedCharacters 使用的 if-else 分支可直接使用 Option.when 简化:
val reservedCharacters =
Option.when(c.hasPath("reserved-characters")) {
c.getString("reserved-characters")
}
目前代码整体看起来较为简洁,最后一步就是验证代码风格是否遵守项目规范了。
其他#
这个问题看起来是一个小问题,但反映出来的现象却可见一斑。首先,对于代码维护者或产品厂商来说要严格遵守 规范,规范才是一切标准的准绳;其次,对于应用设计者而言,良好的技术远见和背景可以提前预见很多潜在风险从而避免一些奇怪的问题,而且架构的防御性设计是衡量软件可扩展性和稳定运行的重要维度。