【问题标题】:Testing Schema Update in Apache Beam在 Apache Beam 中测试模式更新
【发布时间】:2021-10-31 03:11:36
【问题描述】:

我正在创建一个无限制的 Dataflow 管道,并希望确保架构的新版本与旧版本兼容,以便可以在不停止的情况下对其进行更新。如果我的原始对象定义为:

@DefaultSchema(JavaFieldSchema.class)
public class TransactionPojoV1 {
  public final String bank;
  public final double purchaseAmount;

  @SchemaCreate
  public TransactionPojoV1(String bank, double purchaseAmount) {
    this.bank = bank;
    this.purchaseAmount = purchaseAmount;
  }
}

我想添加一个新字段fee

@DefaultSchema(JavaFieldSchema.class)
public class TransactionPojoV2 {
  public final String bank;
  public final double purchaseAmount;
  public final double fee;

  @SchemaCreate
  public TransactionPojoV2(String bank, double purchaseAmount, double fee) {
    this.bank = bank;
    this.purchaseAmount = purchaseAmount;
    this.fee = fee
  }
}

我如何编写一个测试来测试TransactionPojoV2 是否可以从TransactionPojoV1 解码?并确保行为符合预期。

以上内容可能无法通过此测试,不确定,但我想要以下内容:

TransactionPojoV1 transactionPojoV1 = ...

byte[] encoded = Coder.encode(transactionPojoV1);

TransactionPojoV2 transactionPojoV2 = Coder.decode(encoded);

// Assert values are as expected.

我只是不知道该怎么做。

【问题讨论】:

    标签: java google-cloud-dataflow apache-beam


    【解决方案1】:

    您可以从PipelineSchemaRegistry 获取Coder。对于问题中给出的课程,您可以这样做:

    @RunWith(JUnit4.class)
    public class TransactionPojoTest {
        @Rule
        public TestPipeline p = TestPipeline.create();
    
        @Test
        public void testTransactionV1toV2() throws NoSuchSchemaException, IOException {
            SchemaRegistry schemaRegistry = p.getSchemaRegistry();
    
            Coder<TransactionPojoV1> transactionPojoV1Coder = schemaRegistry.getSchemaCoder(TransactionPojoV1.class);
            Coder<TransactionPojoV2> transactionPojoV2Coder = schemaRegistry.getSchemaCoder(TransactionPojoV2.class);
    
            TransactionPojoV1 inputTransaction = new TransactionPojoV1("Some Bank", 100.0);
    
            ByteArrayOutputStream byteArrayOutputStream = new ByteArrayOutputStream();
            transactionPojoV1Coder.encode(inputTransaction, byteArrayOutputStream);
    
            ByteArrayInputStream byteArrayInputStream = new ByteArrayInputStream(byteArrayOutputStream.toByteArray());
    
            TransactionPojoV2 outputTransaction = transactionPojoV2Coder.decode(byteArrayInputStream);
    
            assertEquals(inputTransaction.bank, outputTransaction.bank);
            assertEquals(inputTransaction.purchaseAmount, outputTransaction.purchaseAmount, 1e-15);
            assertEquals(0D, outputTransaction.fee.doubleValue(), 1e-15);
        }
    }
    

    为此,TransactionPojoV2 需要:

    @DefaultSchema(JavaFieldSchema.class)
    public class TransactionPojoV2 {
        public final String bank;
        public final double purchaseAmount;
        @Nullable public final Double fee;
    
        @SchemaCreate
        public TransactionPojoV2(String bank, double purchaseAmount, @Nullable Double fee) {
            this.bank = bank;
            this.purchaseAmount = purchaseAmount;
            this.fee = fee != null ? fee : 0;
        }
    }
    

    【讨论】:

    • 我认为这可能会引发一些误报。当通过重新排序字段以确保架构兼容性来添加新字段时,数据流可以更加灵活,但这里没有考虑到这一点。
    【解决方案2】:

    Dataflow 运行器在这里可以更加宽松,因为它支持重新排序字段以确保更新兼容性。如果有新字段,它会将它们移动到 Schema 的“末尾”,并确保现有字段在新旧模式中的顺序相同。

    这一切对您作为 Beam 用户来说都是透明的,这只是意味着只要您的新架构仅添加相对于旧架构的字段,它就应该是更新兼容的。

    为了测试这种更新兼容性,我建议你使用这样的东西:

    Schema v1 = schemaRegistry.getSchema(TransactionPojoV1.class);
    Schema v2 = schemaRegistry.getSchema(TransactionPojoV2.class);
    
    Set<Field> removedFields = Sets.difference(ImmutableSet.of(v1.getFields()), 
                                               ImmutableSet.of(v2.getFields()));
    assertEmpty(removedFields);
    

    请注意,这只是现成的,我还没有测试过。它还依赖于 guava 中的 Sets 和 ImmutableSet。

    另外,值得注意的是,Dataflow 在实践中可能比这更宽松。 Dataflow 只关心融合边界(即存在 GroupByKey/shuffle 的地方)的更新兼容性。

    【讨论】:

    • 这是否意味着TestPipeline中生成的Schema将与Dataflow上运行的管道的Schema不匹配?由于Schema 包括位置/索引。除非它在编码为字节后进行改组?这意味着我建议的编码/解码测试实际上是无效的,因为它没有测试相同的Schema。这种行为是否记录在任何地方?
    • 或者是Coder不一样?
    • 是的,Schema 将被修改,但它以对用户透明的方式完成。有一个间接层(参见 encoding_position,在 schema.proto 中),它允许运行者设置字段的编码顺序,以确保更新兼容性,同时在用户代码中保持用户定义的顺序。
    猜你喜欢
    • 2020-10-06
    • 2021-06-16
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-02-15
    • 2023-04-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多