MCPcopy Create free account

hub / github.com/delta-io/kafka-delta-ingest / functions

Functions298 in github.com/delta-io/kafka-delta-ingest

Methoddecompress
(bytes: &[u8])
src/serialization.rs:108
Methoddefault
()
src/lib.rs:308
Functiondelta_stats_test
()
src/writer.rs:1156
Functionend_at_initial_offsets
()
tests/offset_tests.rs:343
Methodextract_message_fingerprint
(msg: &[u8])
src/serialization.rs:429
Methodfingerprint_to_i64
(msg: [u8; 8])
src/serialization.rs:437
Methodfmt
(&self, f: &mut fmt::Formatter<'_>)
src/cursor.rs:16
Methodfor_table
Creates a DataWriter to write to the given table
src/writer.rs:340
Methodfrom
(err: InvalidTypeError)
src/transforms.rs:57
Methodfrom
(value: ParquetError)
src/writer.rs:144
Methodfrom_failed_deserialization
Creates a dead letter from bytes that failed deserialization. `json_string` will always be `None`.
src/dead_letters.rs:43
Methodfrom_failed_parquet_row
Creates a dead letter from a record that fails on parquet write. `base64_bytes` will always be `None`. `json_string` will contain the stringified JSON
src/dead_letters.rs:74
Methodfrom_failed_transform
Creates a dead letter from a failed transform. `base64_bytes` will always be `None`.
src/dead_letters.rs:58
Methodfrom_options
( options: DeadLetterQueueOptions, )
src/dead_letters.rs:248
Methodfrom_schema_registry
(sr_settings: SrSettings)
src/serialization.rs:341
Methodfrom_str
(path: &str)
src/transforms.rs:328
Methodfrom_transforms
Creates a new transformer which executes each provided transform on the passed JSON value. Transforms should be provided as a HashMap where the key i
src/transforms.rs:383
Functionget_avro_argument
()
src/main.rs:555
Functionget_end_at_last_offset
()
src/main.rs:592
Functionget_json_argument
()
src/main.rs:530
Functionget_stats_value
(add: &Add, key: &str)
tests/delta_partitions_tests.rs:137
Functionjmespath_epoch_micros_to_iso8601
( args: &[Rcvar], context: &mut Context, )
src/transforms.rs:273
Functionjmespath_epoch_millis_to_iso8601
( args: &[Rcvar], context: &mut Context, )
src/transforms.rs:263
Functionjmespath_epoch_millis_to_micro
( args: &[Rcvar], context: &mut Context, )
src/transforms.rs:282
Functionjmespath_epoch_seconds_to_iso8601
( args: &[Rcvar], context: &mut Context, )
src/transforms.rs:253
Methodjson_default
(decompress_gzip: bool)
src/serialization.rs:70
Functionmain
()
build.rs:5
Functionmain
()
src/main.rs:46
Functionmsg
(s: String)
tests/delta_partitions_tests.rs:142
Methodnew
Creates a new [`ValueBuffer`] to store messages from a Kafka partition.
src/value_buffers.rs:88
Methodnew
Creates a new ingest [`IngestProcessor`].
src/lib.rs:757
Methodnew
Creates an instance of [`IngestMetrics`] for sending metrics to statsd.
src/metrics.rs:30
Methodnew
(context: &Context, expected: &str, actual: String, position: usize)
src/transforms.rs:44
Methodnew
( arrow_schema: Arc<ArrowSchema>, writer_properties: WriterProperties, )
src/writer.rs:297
Methodnew
(decompress_gzip: bool)
src/serialization.rs:104
Methodnew
Create a new SliceableCursor
src/cursor.rs:28
Methodnew
(id: u32, color: &str)
tests/delta_partitions_tests.rs:21
Methodnew
(id: u64)
tests/offset_tests.rs:23
Methodnew
(topic: &str, table: &str, options: IngestOptions)
tests/helpers/mod.rs:573
Methodnew_color_null
(id: u32)
tests/delta_partitions_tests.rs:28
Methodnew_underlying_writer
( cursor: InMemoryWriteableCursor, arrow_schema: Arc<ArrowSchema>, writer_properties:
src/writer.rs:329
Functionparse_fields
(schema: &Value)
tests/helpers/mod.rs:165
Functionparse_json_field
(value: &Value, key: &str)
tests/helpers/mod.rs:549
Functionparse_seek_offsets_test
()
src/main.rs:524
Functionparse_type
(schema: &Value)
tests/helpers/mod.rs:149
Methodpost_rebalance
(&self, _consumer: &BaseConsumer<KafkaContext>, rebalance: &Rebalance)
src/lib.rs:1258
Methodpre_rebalance
(&self, _consumer: &BaseConsumer<KafkaContext>, rebalance: &Rebalance)
src/lib.rs:1237
Functionread_all_slice
()
src/cursor.rs:172
Functionread_all_whole
()
src/cursor.rs:166
Methodread_single_schema_file
( path: &PathBuf, )
src/serialization.rs:395
Methodreset
Clears all value buffers currently held in memory.
src/value_buffers.rs:71
Functionschema_update_test
()
tests/schema_update_tests.rs:25
Functionseek_cursor_current
()
src/cursor.rs:186
Functionseek_cursor_end
()
src/cursor.rs:194
Functionseek_cursor_error_too_long
()
src/cursor.rs:202
Functionseek_cursor_error_too_short
()
src/cursor.rs:212
Functionseek_cursor_start
()
src/cursor.rs:178
Functionsend_kv_json
( producer: &FutureProducer, topic: &str, key: String, json: &Value, )
tests/helpers/mod.rs:92
Functionset_value_sets_recursively
()
src/transforms.rs:528
Functionsubstr_returns_liam_from_william
()
src/transforms.rs:507
Functionsubstr_returns_will_from_william
()
src/transforms.rs:486
Functiontest_avro_default
()
tests/deserialization_tests.rs:164
Functiontest_avro_single_object_encoding_with_file
()
tests/deserialization_tests.rs:296
Functiontest_avro_with_file
()
tests/deserialization_tests.rs:207
Functiontest_avro_with_registry
()
tests/deserialization_tests.rs:252
Functiontest_coercion_tree
()
src/coercions.rs:246
Functiontest_coercions
()
src/coercions.rs:320
Functiontest_delta_partitions
()
tests/delta_partitions_tests.rs:34
Functiontest_dlq
()
tests/dead_letter_tests.rs:29
Functiontest_dont_write_an_empty_buffer
()
tests/buffer_flush_tests.rs:66
Functiontest_epoch_millis_to_micro
()
src/transforms.rs:609
Functiontest_flush_on_size_without_latency_expiration
()
tests/buffer_flush_tests.rs:111
Functiontest_flush_when_latency_expires
()
tests/buffer_flush_tests.rs:14
Functiontest_iso8601_from_epoch_micros_test
()
src/transforms.rs:601
Functiontest_iso8601_from_epoch_seconds_test
()
src/transforms.rs:593
Functiontest_json_default
()
tests/deserialization_tests.rs:39
Functiontest_json_with_args
()
tests/deserialization_tests.rs:79
Functiontest_json_with_registry
()
tests/deserialization_tests.rs:120
Functiontest_start_from_earliest
()
tests/offset_tests.rs:179
Functiontest_start_from_explicit
()
tests/offset_tests.rs:109
Functiontest_start_from_latest
()
tests/offset_tests.rs:239
Functiontest_string_to_timestamp
()
src/coercions.rs:148
Functiontest_transforms_with_epoch_seconds_to_iso8601
()
src/transforms.rs:644
Functiontest_transforms_with_kafka_meta
()
src/transforms.rs:719
Functiontransforms_with_substr
()
src/transforms.rs:554
Methodtry_build
( input_format: &MessageFormat, decompress_gzip: bool, // Add this parameter )
src/serialization.rs:32
Methodtry_from_path
(path: &PathBuf)
src/serialization.rs:369
Methodtry_from_schema_file
(file: &PathBuf)
src/serialization.rs:349
Functionvalue_buffers_conflict_offsets_test
()
src/value_buffers.rs:190
Functionvalue_buffers_test
()
src/value_buffers.rs:132
Methodvec_from_failed_parquet_rows
Creates a vector of tuples where the first element is the stringified JSON value that was not writeable to parquet and the second element is the `Parq
src/dead_letters.rs:90
Functionwhen_both_workers_started_simultaneously
()
tests/emails_s3_tests.rs:23
Functionwhen_both_workers_started_simultaneously_azure
()
tests/emails_azure_blob_tests.rs:19
Functionwhen_rebalance_happens
()
tests/emails_s3_tests.rs:29
Functionwhen_rebalance_happens_azure
()
tests/emails_azure_blob_tests.rs:25
Methodwrite
(&mut self, buf: &[u8])
src/cursor.rs:127
Functionwrite_offsets_to_delta_test
()
src/offsets.rs:171
Functionzero_offset_issue
()
tests/offset_tests.rs:33
← previous201–298 of 298, ranked by callers