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

运行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的凭证类会导致类型匹配失败。

修复步骤

  1. 调整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();
        }
    }
    
  2. 验证依赖版本兼容性
    确保项目依赖的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>
      
  3. 移除V1相关依赖(如果存在)
    检查项目中是否有AWS SDK V1的依赖(如aws-java-sdk-timestreamwrite),如有需排除,避免版本冲突。

内容的提问来源于stack exchange,提问作者Ciaran George

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 23:57:01