运行Apache NiFi Timestream处理器集成测试时遇AWS V2 API兼容错误
问题:NiFi AWS V2 API升级时的Timestream处理器认证错误
错误信息
java.lang.IllegalArgumentException: Invalid option: software.amazon.awssdk.awscore.client.config.AwsClientOption@71168c46. Required value of type interface software.amazon.awssdk.identity.spi.IdentityProvider, but was class org.apache.nifi.processors.aws.credentials.provider.PropertiesCredentialsProvider.
相关代码
集成测试代码
private final String CREDENTIALS_FILE = System.getProperty("user.home") + "/aws-credentials.properties"; private final String DB_NAME = "example_db"; private final String TBL_NAME = "example_tbl"; private final String REGION = Region.EU_WEST_1.id(); @Test public void testSimplePut() throws IOException { final TestRunner runner = TestRunners.newTestRunner(new PutTimestream()); runner.setProperty(PutTimestream.TBL_NAME, TBL_NAME); runner.setProperty(PutTimestream.DB_NAME, DB_NAME); runner.setProperty(PutTimestream.REGION, REGION); runner.setProperty(PutTimestream.CREDENTIALS_FILE, CREDENTIALS_FILE); final Map<String, String> attrs = new HashMap<>(); attrs.put("filename", "timestream_valid.json"); runner.enqueue(Paths.get("src/test/resources/timestream_valid.json"), attrs); runner.run(1); runner.assertAllFlowFilesTransferred(PutTimestream.REL_SUCCESS, 1); }
AbstractTimestreamProcessor代码
public abstract class AbstractTimestreamProcessor extends AbstractAwsProcessor<TimestreamWriteClient, TimestreamWriteClientBuilder> { @Override protected TimestreamWriteClientBuilder createClientBuilder(ProcessContext processContext) { return TimestreamWriteClient.builder(); } }
PutTimestream处理器代码
@Tags( { "amazon", "aws", "timestream", "put" } ) @CapabilityDescription( "AWS Timestream Put Processor." ) public class PutTimestream extends AbstractTimestreamProcessor { public static final PropertyDescriptor DB_NAME = new PropertyDescriptor.Builder() .name("Database Name") .displayName("Database Name") .description("Specifies the name of the Amazon Timestream Database") .required(true) .expressionLanguageSupported(ExpressionLanguageScope.FLOWFILE_ATTRIBUTES) .addValidator(StandardValidators.NON_EMPTY_VALIDATOR) .build(); public static final PropertyDescriptor TBL_NAME = new PropertyDescriptor.Builder() .name("Table Name") .displayName("Table Name") .description("Specifies the name of the Amazon Timestream Table") .required(true) .expressionLanguageSupported(ExpressionLanguageScope.FLOWFILE_ATTRIBUTES) .addValidator(StandardValidators.NON_EMPTY_VALIDATOR) .build(); WriteRecordsResponse writeRecordsResponse; Map<String, String> attributes; public static final List<PropertyDescriptor> properties = Collections .unmodifiableList(Arrays.asList(AWS_CREDENTIALS_PROVIDER_SERVICE, REGION, ACCESS_KEY, SECRET_KEY, CREDENTIALS_FILE, TIMEOUT, DB_NAME, TBL_NAME)); private volatile List<PropertyDescriptor> userDefinedProperties = Collections.emptyList(); private TimestreamRecordConverter timestreamRecordConverter; @Override protected void init(ProcessorInitializationContext context) { timestreamRecordConverter = new TimestreamRecordConverter(getLogger()); } @Override protected List<PropertyDescriptor> getSupportedPropertyDescriptors() { return properties; } private Map<String, String> getAttributes(FlowFile flowFile, WriteRecordsResponse writeRecordsResponse) { attributes = new HashMap(); attributes.putAll(flowFile.getAttributes()); attributes.put(CoreAttributes.MIME_TYPE.key(), "application/json"); attributes.put("timestream.insert.status", String.valueOf(writeRecordsResponse.responseMetadata())); return attributes; } @Override public void onTrigger(ProcessContext context, ProcessSession session) throws ProcessException { FlowFile flowFile = session.get(); if (flowFile == null) { return; } final long startNanos = System.nanoTime(); String dbName = context.getProperty(DB_NAME).evaluateAttributeExpressions(flowFile).getValue(); String tblName = context.getProperty(TBL_NAME).evaluateAttributeExpressions(flowFile).getValue(); final List<Record> records; try { final byte[] content = new byte[(int) flowFile.getSize()]; session.read(flowFile, in -> StreamUtils.fillBuffer(in, content, true)); records = timestreamRecordConverter.convertJsonToTimestreamRecords(new String(content)); WriteRecordsRequest writeRecordsRequest = WriteRecordsRequest.builder() .databaseName(dbName) .tableName(tblName) .records(records) .build(); writeRecordsResponse = getClient(context).writeRecords(writeRecordsRequest); flowFile = session.putAllAttributes(flowFile, getAttributes(flowFile, writeRecordsResponse)); session.transfer(flowFile, REL_SUCCESS); long transmissionMillis = java.util.concurrent.TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startNanos); session.getProvenanceReporter().send(flowFile, tblName, transmissionMillis); } catch (Exception e) { StringWriter sw = new StringWriter(); e.printStackTrace(new PrintWriter(sw)); getLogger().info("Error processing timestream record. Check stacktrace " + sw); session.transfer(flowFile, REL_FAILURE); } } }
解决方案
问题根源
错误本质是凭证提供者类型不兼容:NiFi的PropertiesCredentialsProvider是AWS SDK V1风格的实现,而AWS SDK V2客户端需要的是software.amazon.awssdk.identity.spi.IdentityProvider(或其子类AwsCredentialsProvider),直接传递V1的凭证类会导致类型匹配失败。
修复步骤
调整AbstractTimestreamProcessor,适配V2凭证体系
重写凭证提供者获取逻辑,将NiFi配置的凭证转换为AWS V2兼容的AwsCredentialsProvider实现,同时正确配置客户端构建器:public abstract class AbstractTimestreamProcessor extends AbstractAwsProcessor<TimestreamWriteClient, TimestreamWriteClientBuilder> { @Override protected TimestreamWriteClientBuilder createClientBuilder(ProcessContext processContext) { TimestreamWriteClientBuilder builder = TimestreamWriteClient.builder(); // 配置V2凭证提供者 AwsCredentialsProvider credentialsProvider = getV2CredentialsProvider(processContext); if (credentialsProvider != null) { builder.credentialsProvider(credentialsProvider); } // 配置区域 String region = processContext.getProperty(REGION).getValue(); if (region != null) { builder.region(Region.of(region)); } return builder; } private AwsCredentialsProvider getV2CredentialsProvider(ProcessContext context) { // 处理ACCESS_KEY/SECRET_KEY配置 String accessKey = context.getProperty(ACCESS_KEY).getValue(); String secretKey = context.getProperty(SECRET_KEY).getValue(); if (accessKey != null && secretKey != null) { return StaticCredentialsProvider.create(AwsBasicCredentials.create(accessKey, secretKey)); } // 处理CREDENTIALS_FILE配置 String credentialsFile = context.getProperty(CREDENTIALS_FILE).getValue(); if (credentialsFile != null) { try { Properties props = new Properties(); props.load(new FileInputStream(credentialsFile)); String ak = props.getProperty("accessKey"); String sk = props.getProperty("secretKey"); return StaticCredentialsProvider.create(AwsBasicCredentials.create(ak, sk)); } catch (IOException e) { throw new ProcessException("Failed to load credentials file: " + credentialsFile, e); } } // fallback到默认凭证链(环境变量、实例角色等) return DefaultCredentialsProvider.create(); } }验证依赖版本兼容性
确保项目依赖的NiFi和AWS SDK V2版本匹配:- NiFi版本建议使用1.18.0+(该版本后全面支持AWS SDK V2)
- AWS Timestream Write SDK依赖使用V2版本:
<dependency> <groupId>software.amazon.awssdk</groupId> <artifactId>timestreamwrite</artifactId> <version>2.20.0+</version> </dependency> - NiFi AWS API依赖:
<dependency> <groupId>org.apache.nifi</groupId> <artifactId>nifi-aws-api</artifactId> <version>${nifi.version}</version> </dependency>
移除V1相关依赖(如果存在)
检查项目中是否有AWS SDK V1的依赖(如aws-java-sdk-timestreamwrite),如有需排除,避免版本冲突。
内容的提问来源于stack exchange,提问作者Ciaran George
相关产品推荐
相关产品推荐

