Skip to content

[improve][broker] PIP-246: Improved PROTOBUF_NATIVE schema compatibility checks without using avro-protobuf - #19566

Open
Denovo1998 wants to merge 9 commits into
apache:masterfrom
Denovo1998:protobuf-native-improve
Open

Denovo1998 wants to merge 9 commits into
apache:masterfrom
Denovo1998:protobuf-native-improve

Conversation

@Denovo1998

@Denovo1998 Denovo1998 commented Feb 19, 2023 •

Copy link
Copy Markdown
Contributor

PIP: PIP-246

Depends on #26725 for shared descriptor deserialization fixes: import handling, nested-root resolution, and Java feature preservation. These changes currently remain in this diff and affect the default checker, broker schema validation, and generic client readers, regardless of whether the advanced checker is enabled.

Motivation

The default PROTOBUF_NATIVE checker compares only root message names, so it can accept field changes that break message reading. This PR introduces an opt-in checker that compares Protobuf descriptors directly, without converting them to Avro. The existing checker remains the default, and the schema storage format is unchanged.

Modifications

  • Add ProtobufNativeSchemaAdvancedCompatibilityCheck with directional checks for fields, required values, defaults, maps, oneofs, enums, and supported Protobuf features, including native components of KEY_VALUE schemas.
  • Apply the conservative compatibility policies defined in PIP-246, including rejecting scalar type changes and unsupported features.
  • Bound comparison work per historical schema pair and direction, memoize message and enum pairs, and return COMPARISON_LIMIT_EXCEEDED when the budget is exhausted.
  • Reuse byte-identical schema pairs while still checking other selected historical versions. Include the conflicting stored schema version in bounded error messages, including for key/value components.
  • Fail startup when a configured checker cannot be instantiated. When the advanced checker is selected, also reject conflicting native checkers and unavailable schema storage.
  • Add descriptor, registry, consumer, configuration, and generated-message regressions, plus a compatibility benchmark and appropriate lightproto exclusions.

Activation requires consistent checker configuration and broker restarts. Reconnecting consumers can still be rejected against other retained schema versions, even when their own schema is unchanged. Restore the default checker before downgrading to binaries that do not contain the advanced checker.

Verifying this change

  • Make sure that the change passes the CI checks.

(Please pick either of the following options)

This change is a trivial rework / code cleanup without any test coverage.

(or)

This change is already covered by existing tests, such as (please describe tests).

(or)

This change added tests and can be verified as follows:

(example:)

  • Added integration tests for end-to-end deployment with large payloads (10MB)
  • Extended integration test for recovery after broker failure

Does this pull request potentially affect one of the following parts:

If the box was checked, please highlight the changes

  • Dependencies (add or upgrade a dependency)
  • The public API
  • The schema
  • The default values of configurations
  • The threading model
  • The binary protocol
  • The REST endpoints
  • The admin CLI options
  • The metrics
  • Anything that affects deployment

…tations in the conf; fix test; when the field type changes there is no check, only a warning.
@github-actions

github-actions Bot commented Apr 5, 2023

Copy link
Copy Markdown

The pr had no activity for 30 days, mark with Stale label.

@github-actions github-actions Bot added the Stale label Apr 5, 2023
@Technoboy- Technoboy- added this to the 3.2.0 milestone Jul 31, 2023
@Technoboy- Technoboy- modified the milestones: 3.2.0, 3.3.0 Dec 22, 2023
@coderzc coderzc modified the milestones: 3.3.0, 3.4.0 May 8, 2024
@lhotari lhotari modified the milestones: 4.0.0, 4.1.0 Oct 14, 2024
@coderzc coderzc modified the milestones: 4.1.0, 4.2.0 Sep 1, 2025
@Denovo1998 Denovo1998 closed this Apr 3, 2026
# Conflicts:
#	pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/ProtobufNativeSchemaUtils.java
@Denovo1998 Denovo1998 reopened this Sep 23, 2026
@Denovo1998
Denovo1998 marked this pull request as draft September 23, 2026 12:40
@Denovo1998
Denovo1998 marked this pull request as ready for review September 23, 2026 14:45

@lhotari lhotari left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the continued work on this — the rule engine is careful, and it is mostly right on protobuf semantics: fields matched by number, directional required, the oneof rule including proto3 optional, open/closed enums including the proto2 default-reordering case, packed/unpacked in both directions, and fail-closed handling of unknown editions and features. Keeping the default checker untouched and failing fast when two PROTOBUF_NATIVE checkers are configured are both good.

Two things need resolving before this can merge: the comparison has no work bound and runs synchronously on the registry path (inline on ProtobufNativeSchemaCompatibility.java), and enabling the checker can refuse consumers that reconnect with an unchanged schema under the transitive strategies.

Separately, the description only links to #26695 and leaves the template placeholders. Please summarise the change there, including that the ProtobufNativeSchemaUtils.deserialize() changes reach every native-schema user — the default checker, the upload validator and the client's generic readers — not only those who opt in. Those parts (cycle detection, nested-root resolution without a package, Java feature parsing) could arguably be a separate [fix] PR. The @FieldContext doc of schemaRegistryCompatibilityCheckers (ServiceConfiguration.java:3861) should also get the activation note you added to broker.conf.

return field.legacyEnumFieldTreatedAsClosed();
}

private static void enqueue(ArrayDeque<MessagePair> queue,

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[PERFORMANCE] The pair walk has no work bound: coprime cycles visit n×m pairs, and enum comparison adds F×E per pair

visited deduplicates per (writer, reader) pair, which guarantees termination but not a small bound. With a writer message cycle of length n and a reader cycle of length m, every step compatible, the walk (Wᵢ, Rⱼ) → (Wᵢ₊₁, Rⱼ₊₁) visits lcm(n, m) pairs — n×m when they are coprime. Each pair re-sorts fields (:350-353) and builds bounded path strings. On top of that, compareEnum (:276-309) rebuilds its maps for every enum-typed field, so a message with F fields sharing one E-value enum costs F×E even with a single message pair.

This runs synchronously in the registry's future chain (SchemaRegistryServiceImpl.java:513-517), once per retained version under *_TRANSITIVE and twice for FULL. A first version is stored without a comparison, so whoever can register schemas (auto-update or admin) can set up the shape.

Could you add a budget on total work (pairs plus fields/enum values visited, e.g. proportional to the size of both descriptor graphs) that fails with a dedicated reason, memoize enum pairs, and add a coprime-cycle test or benchmark @Param?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Added a work budget covering pairs, fields, and enum values, plus enum-pair memoization and cached sorting. Exceeding the budget returns COMPARISON_LIMIT_EXCEEDED. I've also added coprime-cycle and shared-enum regression tests.

throw new IncompatibleSchemaException("SCHEMA_RECONSTRUCTION_FAILED: missing schema");
}
Descriptor proposed = null;
for (SchemaData existingData : from) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[COMPATIBILITY] Under transitive strategies, enabling the checker can refuse consumers presenting an unchanged schema

For schema-bearing, non-AUTO_CONSUME subscriptions under BACKWARD_TRANSITIVE/FULL_TRANSITIVE, checkConsumerCompatibility goes to checkCompatibilityWithAll (SchemaRegistryServiceImpl.java:382-388), which, unlike the latest-only path (SchemaRegistryServiceImpl.java:352-357), has no equality shortcut. canRead validates the supported-feature boundary on both sides first, so a schema with an unsupported construct fails even against itself — testAdvancedNativeUnsupportedHistoryEqualityShortcuts (SchemaServiceTest.java:478-490) pins exactly that. History accepted by the root-name-only checker is likewise re-validated on every reconnect after activation.

Is that intended? Treating byte-identical (existing, proposed) pairs as compatible before graph validation would remove the self-rejection, though not incompatibilities with other retained versions. Either way, broker.conf and the PIP should say that activation can refuse reconnecting consumers, not only new registrations.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Byte-identical pairs now skip graph validation, including in transitive checks. Other historical versions are still checked. Both broker.conf and the PIP now warn that reconnecting consumers can still be rejected.

throw incompatible("CARDINALITY_CHANGED", path, reader.getNumber(),
writer.isRepeated() ? "repeated" : "singular", reader.isRepeated() ? "repeated" : "singular");
}
if (writer.getType() != reader.getType()) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[QUESTION] Every type difference is rejected, including lossless widenings

Question rather than a request. canRead is directional, yet any getType() difference fails, so int32→int64 (and uint32→uint64, sint32→sint64) under BACKWARD is rejected even though an old int32 writer read as int64 is lossless. Starting strict is fine, since relaxing later only accepts more. If it is deliberate, could the PIP list the rules that are stricter than protobuf's own guidance (type widening, singular↔repeated, default changes, name reuse under a new number), so users can tell policy rejections from wire incompatibility?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, I'm keeping this strict for now. The PIP now lists these stricter policies explicitly, including widening, singular/repeated changes, defaults, and renumbering, so they aren't confused with wire incompatibility.

ProtobufNativeSchemaCompatibility.canRead(writer, reader);
} catch (IncompatibleSchemaException e) {
String prefix = "strategy=" + strategy + ", direction=" + direction
+ ", existingSchemaSha256=" + fingerprint(existing.getData()) + ", ";

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[MINOR] A transitive failure identifies the conflicting version only by a hash users can't see

existingSchemaSha256 is the SHA-256 of the stored bytes, while pulsar-admin schemas works with versions, so finding the conflicting entry means fetching and hashing each version. The checker interface doesn't receive versions, so the cleanest fix is probably for the registry, which does, to add the version to the exception; a history index would also need a defined order (KeyValueSchemaCompatibilityCheck.java:94-95 reverses the list). Nit: the concatenation on :100-101 leaves a double space before path=.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The registry now adds the actual existingSchemaVersion, including for key/value component failures. I also fixed the extra space before path=.

Comment thread pulsar-broker/build.gradle.kts Outdated
lightproto {
// Test protos that need standard protobuf (GeneratedMessageV3), not lightproto
excludes.addAll("ProtobufSchemaTest.proto", "DataRecord.proto")
excludes.addAll("ProtobufSchemaTest.proto", "DataRecord.proto", "NativeEditionJavaUtf8.proto")

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[NIT] Four of the five new test protos are not excluded from lightproto

Nit: the tests use the protoc-generated outer classes, so this works either way, but lightproto will also generate unused top-level classes for the other four new test .proto files. Excluding them like the existing entries keeps the generated test sources consistent.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done. All five new test protos are now excluded from lightproto generation.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

doc-not-needed Your PR changes do not impact docs Stale

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants