Code
Hub
Workspaces
Following
Trending
Connect
MCP
copy
Create free account
hub
/
github.com/delta-io/kafka-delta-ingest
/ functions
Functions
298 in github.com/delta-io/kafka-delta-ingest
⨍
Functions
298
◇
Types & classes
61
Method
decompress
(bytes: &[u8])
src/serialization.rs:108
Method
default
()
src/lib.rs:308
Function
delta_stats_test
()
src/writer.rs:1156
Function
end_at_initial_offsets
()
tests/offset_tests.rs:343
Method
extract_message_fingerprint
(msg: &[u8])
src/serialization.rs:429
Method
fingerprint_to_i64
(msg: [u8; 8])
src/serialization.rs:437
Method
fmt
(&self, f: &mut fmt::Formatter<'_>)
src/cursor.rs:16
Method
for_table
Creates a DataWriter to write to the given table
src/writer.rs:340
Method
from
(err: InvalidTypeError)
src/transforms.rs:57
Method
from
(value: ParquetError)
src/writer.rs:144
Method
from_failed_deserialization
Creates a dead letter from bytes that failed deserialization. `json_string` will always be `None`.
src/dead_letters.rs:43
Method
from_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
Method
from_failed_transform
Creates a dead letter from a failed transform. `base64_bytes` will always be `None`.
src/dead_letters.rs:58
Method
from_options
( options: DeadLetterQueueOptions, )
src/dead_letters.rs:248
Method
from_schema_registry
(sr_settings: SrSettings)
src/serialization.rs:341
Method
from_str
(path: &str)
src/transforms.rs:328
Method
from_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
Function
get_avro_argument
()
src/main.rs:555
Function
get_end_at_last_offset
()
src/main.rs:592
Function
get_json_argument
()
src/main.rs:530
Function
get_stats_value
(add: &Add, key: &str)
tests/delta_partitions_tests.rs:137
Function
jmespath_epoch_micros_to_iso8601
( args: &[Rcvar], context: &mut Context, )
src/transforms.rs:273
Function
jmespath_epoch_millis_to_iso8601
( args: &[Rcvar], context: &mut Context, )
src/transforms.rs:263
Function
jmespath_epoch_millis_to_micro
( args: &[Rcvar], context: &mut Context, )
src/transforms.rs:282
Function
jmespath_epoch_seconds_to_iso8601
( args: &[Rcvar], context: &mut Context, )
src/transforms.rs:253
Method
json_default
(decompress_gzip: bool)
src/serialization.rs:70
Function
main
()
build.rs:5
Function
main
()
src/main.rs:46
Function
msg
(s: String)
tests/delta_partitions_tests.rs:142
Method
new
Creates a new [`ValueBuffer`] to store messages from a Kafka partition.
src/value_buffers.rs:88
Method
new
Creates a new ingest [`IngestProcessor`].
src/lib.rs:757
Method
new
Creates an instance of [`IngestMetrics`] for sending metrics to statsd.
src/metrics.rs:30
Method
new
(context: &Context, expected: &str, actual: String, position: usize)
src/transforms.rs:44
Method
new
( arrow_schema: Arc<ArrowSchema>, writer_properties: WriterProperties, )
src/writer.rs:297
Method
new
(decompress_gzip: bool)
src/serialization.rs:104
Method
new
Create a new SliceableCursor
src/cursor.rs:28
Method
new
(id: u32, color: &str)
tests/delta_partitions_tests.rs:21
Method
new
(id: u64)
tests/offset_tests.rs:23
Method
new
(topic: &str, table: &str, options: IngestOptions)
tests/helpers/mod.rs:573
Method
new_color_null
(id: u32)
tests/delta_partitions_tests.rs:28
Method
new_underlying_writer
( cursor: InMemoryWriteableCursor, arrow_schema: Arc<ArrowSchema>, writer_properties:
src/writer.rs:329
Function
parse_fields
(schema: &Value)
tests/helpers/mod.rs:165
Function
parse_json_field
(value: &Value, key: &str)
tests/helpers/mod.rs:549
Function
parse_seek_offsets_test
()
src/main.rs:524
Function
parse_type
(schema: &Value)
tests/helpers/mod.rs:149
Method
post_rebalance
(&self, _consumer: &BaseConsumer<KafkaContext>, rebalance: &Rebalance)
src/lib.rs:1258
Method
pre_rebalance
(&self, _consumer: &BaseConsumer<KafkaContext>, rebalance: &Rebalance)
src/lib.rs:1237
Function
read_all_slice
()
src/cursor.rs:172
Function
read_all_whole
()
src/cursor.rs:166
Method
read_single_schema_file
( path: &PathBuf, )
src/serialization.rs:395
Method
reset
Clears all value buffers currently held in memory.
src/value_buffers.rs:71
Function
schema_update_test
()
tests/schema_update_tests.rs:25
Function
seek_cursor_current
()
src/cursor.rs:186
Function
seek_cursor_end
()
src/cursor.rs:194
Function
seek_cursor_error_too_long
()
src/cursor.rs:202
Function
seek_cursor_error_too_short
()
src/cursor.rs:212
Function
seek_cursor_start
()
src/cursor.rs:178
Function
send_kv_json
( producer: &FutureProducer, topic: &str, key: String, json: &Value, )
tests/helpers/mod.rs:92
Function
set_value_sets_recursively
()
src/transforms.rs:528
Function
substr_returns_liam_from_william
()
src/transforms.rs:507
Function
substr_returns_will_from_william
()
src/transforms.rs:486
Function
test_avro_default
()
tests/deserialization_tests.rs:164
Function
test_avro_single_object_encoding_with_file
()
tests/deserialization_tests.rs:296
Function
test_avro_with_file
()
tests/deserialization_tests.rs:207
Function
test_avro_with_registry
()
tests/deserialization_tests.rs:252
Function
test_coercion_tree
()
src/coercions.rs:246
Function
test_coercions
()
src/coercions.rs:320
Function
test_delta_partitions
()
tests/delta_partitions_tests.rs:34
Function
test_dlq
()
tests/dead_letter_tests.rs:29
Function
test_dont_write_an_empty_buffer
()
tests/buffer_flush_tests.rs:66
Function
test_epoch_millis_to_micro
()
src/transforms.rs:609
Function
test_flush_on_size_without_latency_expiration
()
tests/buffer_flush_tests.rs:111
Function
test_flush_when_latency_expires
()
tests/buffer_flush_tests.rs:14
Function
test_iso8601_from_epoch_micros_test
()
src/transforms.rs:601
Function
test_iso8601_from_epoch_seconds_test
()
src/transforms.rs:593
Function
test_json_default
()
tests/deserialization_tests.rs:39
Function
test_json_with_args
()
tests/deserialization_tests.rs:79
Function
test_json_with_registry
()
tests/deserialization_tests.rs:120
Function
test_start_from_earliest
()
tests/offset_tests.rs:179
Function
test_start_from_explicit
()
tests/offset_tests.rs:109
Function
test_start_from_latest
()
tests/offset_tests.rs:239
Function
test_string_to_timestamp
()
src/coercions.rs:148
Function
test_transforms_with_epoch_seconds_to_iso8601
()
src/transforms.rs:644
Function
test_transforms_with_kafka_meta
()
src/transforms.rs:719
Function
transforms_with_substr
()
src/transforms.rs:554
Method
try_build
( input_format: &MessageFormat, decompress_gzip: bool, // Add this parameter )
src/serialization.rs:32
Method
try_from_path
(path: &PathBuf)
src/serialization.rs:369
Method
try_from_schema_file
(file: &PathBuf)
src/serialization.rs:349
Function
value_buffers_conflict_offsets_test
()
src/value_buffers.rs:190
Function
value_buffers_test
()
src/value_buffers.rs:132
Method
vec_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
Function
when_both_workers_started_simultaneously
()
tests/emails_s3_tests.rs:23
Function
when_both_workers_started_simultaneously_azure
()
tests/emails_azure_blob_tests.rs:19
Function
when_rebalance_happens
()
tests/emails_s3_tests.rs:29
Function
when_rebalance_happens_azure
()
tests/emails_azure_blob_tests.rs:25
Method
write
(&mut self, buf: &[u8])
src/cursor.rs:127
Function
write_offsets_to_delta_test
()
src/offsets.rs:171
Function
zero_offset_issue
()
tests/offset_tests.rs:33
← previous
201–298 of 298, ranked by callers