Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
15 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions elasticgraph-indexer/lib/elastic_graph/indexer.rb
Original file line number Diff line number Diff line change
Expand Up @@ -85,6 +85,7 @@ def operation_factory
record_preparer_factory: record_preparer_factory,
logger: datastore_core.logger,
skip_derived_indexing_type_updates: config.skip_derived_indexing_type_updates,
skip_record_validation_percents_by_type: config.skip_record_validation_percents_by_type,
configure_record_validator: nil
)
end
Expand Down
34 changes: 32 additions & 2 deletions elasticgraph-indexer/lib/elastic_graph/indexer/config.rb
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@

module ElasticGraph
class Indexer
class Config < Support::Config.define(:latency_slo_thresholds_by_timestamp_in_ms, :skip_derived_indexing_type_updates, :extension_modules)
class Config < Support::Config.define(:latency_slo_thresholds_by_timestamp_in_ms, :skip_derived_indexing_type_updates, :skip_record_validation_percents_by_type, :extension_modules)
json_schema at: "indexer",
optional: false,
description: "Configuration for indexing operations and metrics used by `elasticgraph-indexer`.",
Expand Down Expand Up @@ -43,15 +43,45 @@ class Config < Support::Config.define(:latency_slo_thresholds_by_timestamp_in_ms
{"WidgetWorkspace" => ["ABC12345678"]}
]
},
skip_record_validation_percents_by_type: {
description: "Map of GraphQL type names to the percentage of records of that type whose per-record " \
"JSON schema validation should be skipped. `0` (or an absent key) validates every record of the " \
"type; `100` skips every record; values in between sample, and may be fractional. The decision is " \
"deterministic per event id (`type:id@vversion`), so the same event makes the same choice on every " \
"retry and on every indexer pod. The event envelope (op, id, type, version, json_schema_version, " \
"latency_timestamps) is always validated, regardless of this setting.\n\n" \
"With a large schema the per-record schema walk consumes a significant share of indexing CPU: every " \
"record is checked against every regex, enum, min/max, format, and abstract-type discriminator " \
"defined for its type. Skipping it trades that check for throughput, which is worthwhile when " \
"backfilling data that was already validated upstream. Leaving a percentage of records validated " \
"keeps a canary in place so schema drift still surfaces.\n\n" \
"Note: skipping validation makes malformed-data detection later and less precise. A malformation " \
"found while building an event's operations is still reported as an isolated event failure, carrying " \
"the message validation itself would have produced. But one found only while serializing an " \
"operation for the datastore (an unparsable rollover index timestamp, or a missing custom routing " \
"field) raises an error that fails the entire batch, including the well-formed events in it. Since " \
"such a batch produces no partial-failure response, the queue redelivers all of its events, and the " \
"malformed record fails them again on each retry until it is drained to the dead letter queue. Leave " \
"this empty for live-traffic ingestion.",
type: "object",
patternProperties: {/^[A-Z]\w*$/.source => {type: "number", minimum: 0, maximum: 100}},
additionalProperties: false,
default: {}, # : untyped
examples: [
{}, # : untyped
{"Widget" => 90, "Component" => 100}
]
},
extension_modules: Support::Config::EXTENSION_MODULE_SCHEMA
}

private

def convert_values(skip_derived_indexing_type_updates:, latency_slo_thresholds_by_timestamp_in_ms:, extension_modules:)
def convert_values(skip_derived_indexing_type_updates:, latency_slo_thresholds_by_timestamp_in_ms:, skip_record_validation_percents_by_type:, extension_modules:)
{
skip_derived_indexing_type_updates: skip_derived_indexing_type_updates.transform_values(&:to_set),
latency_slo_thresholds_by_timestamp_in_ms: latency_slo_thresholds_by_timestamp_in_ms,
skip_record_validation_percents_by_type: skip_record_validation_percents_by_type.transform_values(&:to_f),
extension_modules: SchemaArtifacts::RuntimeMetadata::ExtensionLoader.load_component_extensions(extension_modules)
}
end
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
require "elastic_graph/indexer/record_preparer"
require "elastic_graph/support/json_schema/validator_factory"
require "elastic_graph/support/memoizable_data"
require "zlib"

module ElasticGraph
class Indexer
Expand All @@ -23,6 +24,7 @@ class Factory < Support::MemoizableData.define(
:record_preparer_factory,
:logger,
:skip_derived_indexing_type_updates,
:skip_record_validation_percents_by_type,
:configure_record_validator
)
def build(event)
Expand All @@ -40,15 +42,44 @@ def build(event)
return build_failed_result(event, "event payload", error_message)
end

failed_result = validate_record_returning_failure(event, selected_json_schema_version)
failed_result || BuildResult.success(build_all_operations_for(
event,
record_preparer_factory.for_json_schema_version(selected_json_schema_version)
))
graphql_type_name = event.fetch("type")

if skip_validation?(graphql_type_name, event)
build_success_result_isolating_malformed_records(event, graphql_type_name, selected_json_schema_version)
else
validate_record_returning_failure(event, graphql_type_name, selected_json_schema_version) ||
build_success_result(event, selected_json_schema_version, type_with_skipped_validation: nil)
end
end

private

def build_success_result(event, selected_json_schema_version, type_with_skipped_validation:)
BuildResult.success(
build_all_operations_for(event, record_preparer_factory.for_json_schema_version(selected_json_schema_version)),
type_with_skipped_validation: type_with_skipped_validation
)
end

# Builds the operations for an event whose per-record validation we skipped.
#
# Skipping validation means malformed data the schema walk would have rejected surfaces instead as
# an exception while we build the event's operations, and there is no bounded list of error types to
# enumerate. So we rescue anything and then run the validation we skipped: the validator tells us
# whether the data was actually bad, and if it was, hands the caller the same pinpointed message it
# would have gotten had we validated up front. A clean bill of health from the validator means the
# error was never about the data (a schema artifact defect, or a bug) and must not be swallowed.
def build_success_result_isolating_malformed_records(event, graphql_type_name, selected_json_schema_version)
build_success_result(event, selected_json_schema_version, type_with_skipped_validation: graphql_type_name)
rescue => exception
failed_result = validate_record_returning_failure(event, graphql_type_name, selected_json_schema_version)
# `raise` is overridden below to stop this class from *originating* an error instead of returning a
# `BuildResult`. Here we propagate one that already escaped a collaborator, which is exactly what
# happens without this rescue, so we deliberately bypass that guard.
::Kernel.raise(exception) unless failed_result
failed_result
end

def select_json_schema_version(event)
available_json_schema_versions = schema_artifacts.available_json_schema_versions

Expand Down Expand Up @@ -117,23 +148,57 @@ def prepare_event(event)
event.merge("record" => event["record"].merge("id" => event.fetch("id")))
end

def validate_record_returning_failure(event, selected_json_schema_version)
def validate_record_returning_failure(event, graphql_type_name, selected_json_schema_version)
record = event.fetch("record")
graphql_type_name = event.fetch("type")
validator = validator(graphql_type_name, selected_json_schema_version)

if (error_message = validator.validate_with_error_message(record))
build_failed_result(event, "#{graphql_type_name} record", error_message)
end
end

# `Zlib.crc32` returns a value in `[0, 2**32)`. Pre-dividing that space by 100 lets us test a
# configured percent with a single multiply instead of dividing on every event.
CRC32_SPACE_PER_PERCENT = (1 << 32) / 100.0

# Decides whether to skip per-record validation for `event` of `type`. The decision is
# deterministic per event id: a stable `Zlib.crc32` of `EventID#to_s` maps each event to a
# point in the CRC32 space, and we skip validation for the configured percentage of that
# space. Same event id => same decision across pods and retries, so retries never flip a
# record between validated and skipped. `String#hash` is unsuitable here, as `RUBY_HASH_SEED`
# is per-process. The `<= 0` and `>= 100` guards keep the endpoints exact, so no float
# boundary error can make a `0` percent skip a record or a `100` percent validate one.
def skip_validation?(type, event)
percent = skip_record_validation_percents_by_type[type]
return false if percent.nil? || percent <= 0
return true if percent >= 100
::Zlib.crc32(EventID.from_event(event).to_s) < percent * CRC32_SPACE_PER_PERCENT
end

def build_failed_result(event, payload_description, validation_message)
message = "Malformed #{payload_description}. #{validation_message}"

# Here we use the `RecordPreparer::Identity` record preparer because we may not have a valid JSON schema
# version number in this case (which is usually required to get a `RecordPreparer` from the factory), and
# we won't wind up using the record preparer for real on these operations, anyway.
operations = build_all_operations_for(event, RecordPreparer::Identity)
#
# Building operations for an event we already know is malformed can itself fail--for example, when the
# record omits a field an update target derives its id from. Reporting what was malformed matters more
# than reporting the operations we would have run, and `FailedEventError#operations` is documented to
# sometimes be empty for exactly this reason, so we fall back to no operations rather than let a second
# failure mask the first.
operations = begin
build_all_operations_for(event, RecordPreparer::Identity)
rescue => exception
Comment thread
vermatron marked this conversation as resolved.
logger.warn({
"message_type" => "FailedEventOperationBuildingFailure",
"message_id" => event["message_id"],
"event_id" => EventID.from_event(event).to_s,
"error_class" => exception.class.name,
"error_message" => exception.message
})
[] # : ::Array[_Operation]
end

BuildResult.failure(FailedEventError.new(event: event, operations: operations.to_set, main_message: message))
end
Expand Down Expand Up @@ -192,14 +257,17 @@ def raise(*args)
# Return value from `build` that indicates what happened.
# - If it was successful, `operations` will be a non-empty array of operations and `failed_event_error` will be nil.
# - If there was a validation issue, `operations` will be an empty array and `failed_event_error` will be non-nil.
BuildResult = ::Data.define(:operations, :failed_event_error) do
# - `type_with_skipped_validation` names the event's GraphQL type when per-record validation was skipped
# (via `skip_record_validation_percents_by_type`), and is nil otherwise. `Processor` aggregates this
# for observability.
BuildResult = ::Data.define(:operations, :failed_event_error, :type_with_skipped_validation) do
# @implements BuildResult
def self.success(operations)
new(operations, nil)
def self.success(operations, type_with_skipped_validation: nil)
new(operations, nil, type_with_skipped_validation)
end

def self.failure(failed_event_error)
new([], failed_event_error)
new([], failed_event_error, nil)
end
end
end
Expand Down
22 changes: 22 additions & 0 deletions elasticgraph-indexer/lib/elastic_graph/indexer/processor.rb
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,8 @@ def process_returning_failures(events, refresh_indices: false)

factory_results = factory_results_by_event.values

log_skipped_record_validations(factory_results)

bulk_result = @datastore_router.bulk(factory_results.flat_map(&:operations), refresh: refresh_indices)
successful_operations = bulk_result.successful_operations(check_failures: false)

Expand All @@ -62,6 +64,26 @@ def process_returning_failures(events, refresh_indices: false)

private

# Emits a single aggregate log line per batch when any records had their per-record validation
# skipped (via `skip_record_validation_percents_by_type`). Skipping a safety check should never be
# silent, but per-record logging would be untenable at backfill scale (a `100` percent is one line
# per record), so we tally by type and log once. Mirrors the batch-level
# `ElasticGraphIndexingLatencies` log.
#
# A skipped record that then failed is deliberately absent from the tally: `Operation::Factory`
# re-runs the skipped validation when building its operations raises, so such a record ended up
# validated after all, and is already reported as a failed event.
def log_skipped_record_validations(factory_results)
counts_by_type = factory_results.filter_map(&:type_with_skipped_validation).tally
return if counts_by_type.empty?

@logger.info({
"message_type" => "RecordValidationSkipped",
"count" => counts_by_type.values.sum,
"counts_by_type" => counts_by_type
})
end

def categorize_failures(failures, events)
source_event_versions_by_cluster_by_op = @datastore_router.source_event_versions_in_index(
failures.flat_map { |f| f.versioned_operations.to_a }
Expand Down
4 changes: 4 additions & 0 deletions elasticgraph-indexer/sig/elastic_graph/indexer/config.rbs
Original file line number Diff line number Diff line change
Expand Up @@ -5,16 +5,19 @@ module ElasticGraph

attr_reader latency_slo_thresholds_by_timestamp_in_ms: ::Hash[::String, ::Integer]
attr_reader skip_derived_indexing_type_updates: ::Hash[::String, ::Set[::String]]
attr_reader skip_record_validation_percents_by_type: ::Hash[::String, ::Float]
attr_reader extension_modules: ::Array[::Module]

def initialize: (
?latency_slo_thresholds_by_timestamp_in_ms: ::Hash[::String, ::Integer],
?skip_derived_indexing_type_updates: ::Hash[::String, ::Set[::String]],
?skip_record_validation_percents_by_type: ::Hash[::String, ::Float],
?extension_modules: ::Array[::Module]) -> void

def with: (
?latency_slo_thresholds_by_timestamp_in_ms: ::Hash[::String, ::Integer],
?skip_derived_indexing_type_updates: ::Hash[::String, ::Set[::String]],
?skip_record_validation_percents_by_type: ::Hash[::String, ::Float],
?extension_modules: ::Array[::Module]) -> Config

def self.members: () -> ::Array[::Symbol]
Expand All @@ -26,6 +29,7 @@ module ElasticGraph
def convert_values: (
latency_slo_thresholds_by_timestamp_in_ms: untyped,
skip_derived_indexing_type_updates: untyped,
skip_record_validation_percents_by_type: untyped,
extension_modules: untyped
) -> ::Hash[::Symbol, untyped]

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ module ElasticGraph
attr_reader record_preparer_factory: RecordPreparer::Factory
attr_reader logger: ::Logger
attr_reader skip_derived_indexing_type_updates: ::Hash[::String, ::Set[::String]]
attr_reader skip_record_validation_percents_by_type: ::Hash[::String, ::Float]
attr_reader configure_record_validator: (^(validatorFactory) -> validatorFactory)?

def initialize: (
Expand All @@ -19,6 +20,7 @@ module ElasticGraph
record_preparer_factory: RecordPreparer::Factory,
logger: ::Logger,
skip_derived_indexing_type_updates: ::Hash[::String, ::Set[::String]],
skip_record_validation_percents_by_type: ::Hash[::String, ::Float],
configure_record_validator: (^(validatorFactory) -> validatorFactory)?
) -> void

Expand All @@ -28,6 +30,7 @@ module ElasticGraph
?record_preparer_factory: RecordPreparer::Factory,
?logger: ::Logger,
?skip_derived_indexing_type_updates: ::Hash[::String, ::Set[::String]],
?skip_record_validation_percents_by_type: ::Hash[::String, ::Float],
?configure_record_validator: (^(validatorFactory) -> validatorFactory)?
) -> instance
end
Expand All @@ -44,7 +47,11 @@ module ElasticGraph

def select_json_schema_version: (event) { (BuildResult) -> bot } -> (::Integer | bot)
def prepare_event: (event) -> event
def validate_record_returning_failure: (event, ::Integer) -> BuildResult?
def validate_record_returning_failure: (event, ::String, ::Integer) -> BuildResult?
def build_success_result: (event, ::Integer, type_with_skipped_validation: ::String?) -> BuildResult
def build_success_result_isolating_malformed_records: (event, ::String, ::Integer) -> BuildResult
CRC32_SPACE_PER_PERCENT: ::Float
def skip_validation?: (::String, event) -> bool
def build_failed_result: (event, ::String, ::String) -> BuildResult
def build_all_operations_for: (event, _RecordPreparer) -> ::Array[_Operation]
def index_definitions_for: (::String) -> ::Array[DatastoreCore::_IndexDefinition]
Expand All @@ -53,15 +60,17 @@ module ElasticGraph
class BuildResult
attr_reader operations: ::Array[_Operation]
attr_reader failed_event_error: FailedEventError?
attr_reader type_with_skipped_validation: ::String?

def initialize: (::Array[_Operation], FailedEventError?) -> void
def initialize: (::Array[_Operation], FailedEventError?, ::String?) -> void

def with: (
?operations: ::Array[_Operation],
?failed_event_error: FailedEventError?
?failed_event_error: FailedEventError?,
?type_with_skipped_validation: ::String?
) -> BuildResult

def self.success: (::Array[_Operation]) -> BuildResult
def self.success: (::Array[_Operation], ?type_with_skipped_validation: ::String?) -> BuildResult
def self.failure: (FailedEventError) -> BuildResult
end
end
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ module ElasticGraph
@indexing_latency_slo_thresholds_by_timestamp_in_ms: ::Hash[::String, ::Integer]
@clock: singleton(::Time)

def log_skipped_record_validations: (::Array[Operation::Factory::BuildResult]) -> void
def categorize_failures: (::Array[FailedEventError], ::Array[event]) -> ::Array[FailedEventError]
def calculate_latency_metrics: (::Array[_Operation], ::Array[Operation::Result]) -> void
end
Expand Down
Loading