This file contains descriptions of all specs and configs used for YTsaurus Flow configuration. The file generation algorithm recursively searches for all subconfigs. It is not perfect.
NYT::NApi::EConnectionType
|
Possible values
|
Description
|
|
native
|
|
|
rpc
|
|
NYT::NApi::NRpcProxy::EAddressType
|
Possible values
|
Description
|
|
internal_rpc
|
|
|
monitoring_http
|
|
|
http
|
|
|
https
|
|
|
public_rpc
|
|
|
chyt_http
|
|
|
chyt_https
|
|
NYT::NApi::NRpcProxy::TConnectionConfig
Source: yt/yt/client/api/rpc_proxy/config.h
|
Parameter
|
Description
|
|
connection_type
|
Type: NYT::NApi::EConnectionType
Default value: native
|
|
cluster_name
|
Type: std::optional<std::string>
|
|
table_mount_cache
|
Type: NYT::TIntrusivePtr<NYT::NApi::TTableMountCacheConfig>
Default value: {}
|
|
replication_card_cache
|
Type: NYT::TIntrusivePtr<NYT::NChaosClient::TReplicationCardCacheConfig>
|
|
chaos_lease_cache
|
Type: NYT::TIntrusivePtr<NYT::NChaosClient::TChaosLeaseCacheConfig>
|
|
cluster_url
|
Type: std::optional<std::string>
|
|
cluster_tag
|
Type: std::optional<NYT::TStrongTypedef<unsigned short, NYT::NObjectClient::TCellTagTag, NYT::TStrongTypedefOptions{true}>>
|
|
proxy_role
|
Type: std::optional<std::string>
|
|
proxy_address_type
|
Type: std::optional<NYT::NApi::NRpcProxy::EAddressType>
|
|
proxy_network_name
|
Type: std::optional<std::string>
|
|
proxy_addresses
|
Type: std::optional<std::vector<std::string>>
|
|
proxy_endpoints
|
Type: NYT::TIntrusivePtr<NYT::NRpc::TServiceDiscoveryEndpointsConfig>
|
|
proxy_unix_domain_socket
|
Type: std::optional<std::string>
|
|
enable_proxy_discovery
|
Type: bool
Default value: true
|
|
proxy_url_aliasing_rules
|
Type: THashMap<std::string, std::string>
Default value: {}
|
|
dynamic_channel_pool
|
Type: NYT::TIntrusivePtr<NYT::NRpc::TDynamicChannelPoolConfig>
Default value: {}
|
|
ping_period
|
Type: TDuration
Default value: 3s
|
|
proxy_list_update_period
|
Type: TDuration
Default value: 5m
|
|
proxy_list_retry_period
|
Type: TDuration
Default value: 1s
|
|
max_proxy_list_retry_period
|
Type: TDuration
Default value: 30s
|
|
max_proxy_list_update_attempts
|
Type: int
Default value: 3
|
|
rpc_timeout
|
Type: TDuration
Default value: 30s
|
|
rpc_acknowledgement_timeout
|
Type: std::optional<TDuration>
Default value: 15000
|
|
timestamp_provider_latest_timestamp_update_period
|
Type: TDuration
Default value: 3s
|
|
default_transaction_timeout
|
Type: TDuration
Default value: 30s
|
|
default_lookup_rows_timeout
|
Type: TDuration
Default value: 30s
|
|
default_select_rows_timeout
|
Type: TDuration
Default value: 30s
|
|
default_total_streaming_timeout
|
Type: TDuration
Default value: 15m
|
|
default_streaming_stall_timeout
|
Type: TDuration
Default value: 1m
|
|
use_total_streaming_timeout_for_heavy_reads
|
Type: bool
Default value: true
|
|
default_chaos_lease_timeout
|
Type: TDuration
Default value: 30s
|
|
default_ping_period
|
Type: TDuration
Default value: 5s
|
|
bus_client
|
Type: NYT::TIntrusivePtr<NYT::NBus::NTcp::TBusConfig>
Default value: {}
|
|
idle_channel_ttl
|
Type: TDuration
Default value: 5m
|
|
http_client
|
Type: NYT::TIntrusivePtr<NYT::NHttp::TClientConfig>
Default value: {}
|
|
https_client
|
Type: NYT::TIntrusivePtr<NYT::NHttps::TClientConfig>
Default value: {}
|
|
request_codec
|
Type: NYT::NCompression::ECodec
Default value: none
|
|
response_codec
|
Type: NYT::NCompression::ECodec
Default value: none
|
|
enable_legacy_rpc_codecs
|
Type: bool
Default value: false
|
|
enable_retries
|
Type: bool
Default value: false
|
|
retrying_channel
|
Type: NYT::TIntrusivePtr<NYT::NRpc::TRetryingChannelConfig>
Default value: {}
|
|
modify_rows_batch_capacity
|
Type: long
Default value: 0
|
|
clock_cluster_tag
|
Type: NYT::TStrongTypedef<unsigned short, NYT::NObjectClient::TCellTagTag, NYT::TStrongTypedefOptions{true}>
Default value: 61444
|
|
udf_registry_path
|
Type: std::optional<TString>
|
|
enable_select_query_tracing_tag
|
Type: bool
Default value: false
|
|
do_not_drop_pure_exclusive_locks
|
Type: bool
Default value: true
|
|
enable_control_multiplexing_band
|
Type: bool
Default value: false
|
NYT::NApi::TTableMountCacheConfig
Source: yt/yt/client/api/config.h
|
Parameter
|
Description
|
|
shard_count
|
Type: unsigned long
Default value: 1
|
|
expire_after_access_time
|
Type: TDuration
Default value: 30m
|
|
expire_after_successful_update_time
|
Type: TDuration
Default value: 30m
|
|
expire_after_failed_update_time
|
Type: TDuration
Default value: 15s
|
|
refresh_time
|
Type: std::optional<TDuration>
Default value: 30000
|
|
expiration_period
|
Type: std::optional<TDuration>
Default value: 10000
|
|
batch_update
|
Type: bool
Default value: false
|
|
reject_if_entry_is_requested_but_not_ready
|
Type: bool
Default value: false
|
|
on_error_retry_count
|
Type: int
Default value: 5
|
|
on_error_retry_slack_period
|
Type: TDuration
Default value: 1s
|
NYT::NBus::EEncryptionMode
|
Possible values
|
Description
|
|
disabled
|
|
|
optional
|
|
|
required
|
|
NYT::NBus::EVerificationMode
|
Possible values
|
Description
|
|
none
|
|
|
ca
|
|
|
full
|
|
NYT::NBus::NTcp::TBusConfig
Source: yt/yt/core/bus/tcp/config.h
|
Parameter
|
Description
|
|
enable_no_delay
|
Type: bool
Default value: true
|
|
enable_aggressive_reconnect
|
Type: bool
Default value: false
|
|
allow_bypass_tls
|
Type: bool
Default value: false
|
|
min_rto
|
Type: TDuration
Default value: 100ms
|
|
max_rto
|
Type: TDuration
Default value: 30s
|
|
rto_scale
|
Type: double
Default value: 2.0
|
|
connect_timeout
|
Type: TDuration
Default value: 15s
|
|
ca
|
Type: NYT::TIntrusivePtr<NYT::NCrypto::TPemBlobConfig>
|
|
cert_chain
|
Type: NYT::TIntrusivePtr<NYT::NCrypto::TPemBlobConfig>
|
|
private_key
|
Type: NYT::TIntrusivePtr<NYT::NCrypto::TPemBlobConfig>
|
|
ssl_configuration_commands
|
Type: std::vector<NYT::TIntrusivePtr<NYT::NCrypto::TSslContextCommand>>
Default value: []
|
|
insecure_skip_verify
|
Type: bool
Default value: false
|
|
enable_quick_ack
|
Type: bool
Default value: true
|
|
bind_retry_count
|
Type: int
Default value: 5
|
|
bind_retry_backoff
|
Type: TDuration
Default value: 3s
|
|
connection_start_delay
|
Type: std::optional<TDuration>
|
|
packet_decoder_delay
|
Type: std::optional<TDuration>
|
|
read_stall_timeout
|
Type: TDuration
Default value: 1m
|
|
write_stall_timeout
|
Type: TDuration
Default value: 1m
|
|
verify_checksums
|
Type: bool
Default value: true
|
|
generate_checksums
|
Type: bool
Default value: true
|
|
enable_local_bypass
|
Type: bool
Default value: true
|
|
encryption_mode
|
Type: NYT::NBus::EEncryptionMode
Default value: optional
|
|
verification_mode
|
Type: NYT::NBus::EVerificationMode
Default value: none
|
|
cipher_list
|
Type: std::optional<std::string>
|
|
load_certs_from_bus_certs_directory
|
Type: bool
Default value: false
|
|
peer_alternative_host_name
|
Type: std::optional<std::string>
|
NYT::NBus::NTcp::TBusServerConfig
Source: yt/yt/core/bus/tcp/config.h
|
Parameter
|
Description
|
|
enable_no_delay
|
Type: bool
Default value: true
|
|
enable_aggressive_reconnect
|
Type: bool
Default value: false
|
|
allow_bypass_tls
|
Type: bool
Default value: false
|
|
min_rto
|
Type: TDuration
Default value: 100ms
|
|
max_rto
|
Type: TDuration
Default value: 30s
|
|
rto_scale
|
Type: double
Default value: 2.0
|
|
connect_timeout
|
Type: TDuration
Default value: 15s
|
|
ca
|
Type: NYT::TIntrusivePtr<NYT::NCrypto::TPemBlobConfig>
|
|
cert_chain
|
Type: NYT::TIntrusivePtr<NYT::NCrypto::TPemBlobConfig>
|
|
private_key
|
Type: NYT::TIntrusivePtr<NYT::NCrypto::TPemBlobConfig>
|
|
ssl_configuration_commands
|
Type: std::vector<NYT::TIntrusivePtr<NYT::NCrypto::TSslContextCommand>>
Default value: []
|
|
insecure_skip_verify
|
Type: bool
Default value: false
|
|
enable_quick_ack
|
Type: bool
Default value: true
|
|
bind_retry_count
|
Type: int
Default value: 5
|
|
bind_retry_backoff
|
Type: TDuration
Default value: 3s
|
|
connection_start_delay
|
Type: std::optional<TDuration>
|
|
packet_decoder_delay
|
Type: std::optional<TDuration>
|
|
read_stall_timeout
|
Type: TDuration
Default value: 1m
|
|
write_stall_timeout
|
Type: TDuration
Default value: 1m
|
|
verify_checksums
|
Type: bool
Default value: true
|
|
generate_checksums
|
Type: bool
Default value: true
|
|
enable_local_bypass
|
Type: bool
Default value: true
|
|
encryption_mode
|
Type: NYT::NBus::EEncryptionMode
Default value: optional
|
|
verification_mode
|
Type: NYT::NBus::EVerificationMode
Default value: none
|
|
cipher_list
|
Type: std::optional<std::string>
|
|
load_certs_from_bus_certs_directory
|
Type: bool
Default value: false
|
|
peer_alternative_host_name
|
Type: std::optional<std::string>
|
|
port
|
Type: std::optional<int>
|
|
unix_domain_socket_path
|
Type: std::optional<std::string>
|
|
max_backlog_size
|
Type: int
Default value: 8192
|
|
max_simultaneous_connections
|
Type: int
Default value: 50000
|
NYT::NBus::NTcp::TDispatcherConfig
Source: yt/yt/core/bus/tcp/config.h
|
Parameter
|
Description
|
|
thread_pool_size
|
Type: int
Default value: 8
|
|
thread_pool_polling_period
|
Type: TDuration
Default value: 10ms
|
|
network_bandwidth
|
Type: std::optional<long>
|
|
networks
|
Type: THashMap<std::string, std::vector<NYT::NNet::TIP6Network>>
Default value: {}
|
|
multiplexing_bands
|
Type: NYT::TEnumIndexedArray<NYT::NBus::EMultiplexingBand, NYT::TIntrusivePtr<NYT::NBus::NTcp::TMultiplexingBandConfig>>
Default value: {}
|
|
bus_certs_directory_path
|
Type: std::optional<std::string>
|
|
enable_local_bypass
|
Type: bool
Default value: false
|
NYT::NBus::NTcp::TDispatcherDynamicConfig
Source: yt/yt/core/bus/tcp/config.h
|
Parameter
|
Description
|
|
thread_pool_size
|
Type: std::optional<int>
|
|
thread_pool_polling_period
|
Type: std::optional<TDuration>
|
|
network_bandwidth
|
Type: std::optional<long>
|
|
networks
|
Type: std::optional<THashMap<std::string, std::vector<NYT::NNet::TIP6Network>>>
|
|
multiplexing_bands
|
Type: std::optional<NYT::TEnumIndexedArray<NYT::NBus::EMultiplexingBand, NYT::TIntrusivePtr<NYT::NBus::NTcp::TMultiplexingBandConfig>>>
|
|
bus_certs_directory_path
|
Type: std::optional<std::string>
|
|
enable_local_bypass
|
Type: std::optional<bool>
|
NYT::NBus::NTcp::TMultiplexingBandConfig
Source: yt/yt/core/bus/tcp/config.h
|
Parameter
|
Description
|
|
tos_level
|
Type: int
Default value: 0
|
|
network_to_tos_level
|
Type: THashMap<std::string, int>
Default value: {}
|
|
min_multiplexing_parallelism
|
Type: int
Default value: 1
|
|
max_multiplexing_parallelism
|
Type: int
Default value: 1000
|
NYT::NChaosClient::TChaosLeaseCacheConfig
Source: yt/yt/client/chaos_client/config.h
|
Parameter
|
Description
|
|
shard_count
|
Type: unsigned long
Default value: 1
|
|
expire_after_access_time
|
Type: TDuration
Default value: 1d
|
|
expire_after_successful_update_time
|
Type: TDuration
Default value: 1d
|
|
expire_after_failed_update_time
|
Type: TDuration
Default value: 1m
|
|
refresh_time
|
Type: std::optional<TDuration>
Default value: YsonEntity
|
|
expiration_period
|
Type: std::optional<TDuration>
Default value: 10000
|
|
batch_update
|
Type: bool
Default value: false
|
|
retry_backoff_time
|
Type: TDuration
Default value: 3s
|
|
retry_attempts
|
Type: int
Default value: 10
|
|
enable_exponential_retry_backoffs
|
Type: bool
Default value: false
|
|
retry_backoff
|
Type: NYT::TExponentialBackoffOptions
Default value:
{
"backoff_jitter" = 0.1;
"backoff_multiplier" = 1.5;
"invocation_count" = 10;
"max_backoff" = 5000;
"min_backoff" = 1000;
}
|
|
retry_timeout
|
Type: std::optional<TDuration>
|
|
discover_timeout
|
Type: TDuration
Default value: 15s
|
|
acknowledgement_timeout
|
Type: TDuration
Default value: 15s
|
|
rediscover_period
|
Type: TDuration
Default value: 1m
|
|
rediscover_splay
|
Type: TDuration
Default value: 15s
|
|
hard_backoff_time
|
Type: TDuration
Default value: 1m
|
|
soft_backoff_time
|
Type: TDuration
Default value: 15s
|
|
max_peer_count
|
Type: int
Default value: 100
|
|
hashes_per_peer
|
Type: int
Default value: 10
|
|
min_peer_count_for_priority_awareness
|
Type: int
Default value: 0
|
|
enable_power_of_two_choices_strategy
|
Type: bool
Default value: true
|
|
max_concurrent_discover_requests
|
Type: int
Default value: 10
|
|
random_peer_eviction_period
|
Type: TDuration
Default value: 1s
|
|
enable_peer_polling
|
Type: bool
Default value: false
|
|
peer_polling_period
|
Type: TDuration
Default value: 1m
|
|
peer_polling_period_splay
|
Type: TDuration
Default value: 10s
|
|
peer_polling_request_timeout
|
Type: TDuration
Default value: 15s
|
|
peer_priority_strategy
|
Type: NYT::NRpc::EPeerPriorityStrategy
Default value: none
|
|
discovery_session_timeout
|
Type: TDuration
|
|
disable_balancing_on_single_address
|
Type: bool
Default value: true
|
|
hedging_delay
|
Type: std::optional<TDuration>
|
|
cancel_primary_request_on_hedging
|
Type: bool
Default value: false
|
|
addresses
|
Type: std::optional<std::vector<std::string>>
|
|
endpoints
|
Type: NYT::TIntrusivePtr<NYT::NRpc::TServiceDiscoveryEndpointsConfig>
|
|
enable_watching
|
Type: bool
|
NYT::NChaosClient::TReplicationCardCacheConfig
Source: yt/yt/client/chaos_client/config.h
|
Parameter
|
Description
|
|
shard_count
|
Type: unsigned long
Default value: 1
|
|
expire_after_access_time
|
Type: TDuration
Default value: 1m
|
|
expire_after_successful_update_time
|
Type: TDuration
Default value: 1m
|
|
expire_after_failed_update_time
|
Type: TDuration
Default value: 1m
|
|
refresh_time
|
Type: std::optional<TDuration>
Default value: 10000
|
|
expiration_period
|
Type: std::optional<TDuration>
Default value: 10000
|
|
batch_update
|
Type: bool
Default value: false
|
|
retry_backoff_time
|
Type: TDuration
Default value: 3s
|
|
retry_attempts
|
Type: int
Default value: 10
|
|
enable_exponential_retry_backoffs
|
Type: bool
Default value: false
|
|
retry_backoff
|
Type: NYT::TExponentialBackoffOptions
Default value:
{
"backoff_jitter" = 0.1;
"backoff_multiplier" = 1.5;
"invocation_count" = 10;
"max_backoff" = 5000;
"min_backoff" = 1000;
}
|
|
retry_timeout
|
Type: std::optional<TDuration>
|
|
discover_timeout
|
Type: TDuration
Default value: 15s
|
|
acknowledgement_timeout
|
Type: TDuration
Default value: 15s
|
|
rediscover_period
|
Type: TDuration
Default value: 1m
|
|
rediscover_splay
|
Type: TDuration
Default value: 15s
|
|
hard_backoff_time
|
Type: TDuration
Default value: 1m
|
|
soft_backoff_time
|
Type: TDuration
Default value: 15s
|
|
max_peer_count
|
Type: int
Default value: 100
|
|
hashes_per_peer
|
Type: int
Default value: 10
|
|
min_peer_count_for_priority_awareness
|
Type: int
Default value: 0
|
|
enable_power_of_two_choices_strategy
|
Type: bool
Default value: true
|
|
max_concurrent_discover_requests
|
Type: int
Default value: 10
|
|
random_peer_eviction_period
|
Type: TDuration
Default value: 1s
|
|
enable_peer_polling
|
Type: bool
Default value: false
|
|
peer_polling_period
|
Type: TDuration
Default value: 1m
|
|
peer_polling_period_splay
|
Type: TDuration
Default value: 10s
|
|
peer_polling_request_timeout
|
Type: TDuration
Default value: 15s
|
|
peer_priority_strategy
|
Type: NYT::NRpc::EPeerPriorityStrategy
Default value: none
|
|
discovery_session_timeout
|
Type: TDuration
|
|
disable_balancing_on_single_address
|
Type: bool
Default value: true
|
|
hedging_delay
|
Type: std::optional<TDuration>
|
|
cancel_primary_request_on_hedging
|
Type: bool
Default value: false
|
|
addresses
|
Type: std::optional<std::vector<std::string>>
|
|
endpoints
|
Type: NYT::TIntrusivePtr<NYT::NRpc::TServiceDiscoveryEndpointsConfig>
|
|
enable_watching
|
Type: bool
|
|
watched_cache
|
Type: NYT::TIntrusivePtr<NYT::NChaosClient::TWatchedReplicationCardCacheConfig>
Default value: {}
|
NYT::NChaosClient::TWatchedReplicationCardCacheConfig
Source: yt/yt/client/chaos_client/config.h
|
Parameter
|
Description
|
|
shard_count
|
Type: unsigned long
Default value: 1
|
|
expire_after_access_time
|
Type: TDuration
Default value: 1d
|
|
expire_after_successful_update_time
|
Type: TDuration
Default value: 1d
|
|
expire_after_failed_update_time
|
Type: TDuration
Default value: 1m
|
|
refresh_time
|
Type: std::optional<TDuration>
Default value: YsonEntity
|
|
expiration_period
|
Type: std::optional<TDuration>
Default value: 10000
|
|
batch_update
|
Type: bool
Default value: false
|
NYT::NClient::NCache::TClientsCacheConfig
Source: yt/yt/client/cache/config.h
NYT::NCodegen::EOptimizationLevel
|
Possible values
|
Description
|
|
none
|
|
|
default
|
|
NYT::NCompression::ECodec
|
Possible values
|
Description
|
|
none
|
|
|
snappy
|
|
|
lz4
|
|
|
lz4_high_compression
|
|
|
brotli_1
|
|
|
brotli_2
|
|
|
brotli_3
|
|
|
brotli_4
|
|
|
brotli_5
|
|
|
brotli_6
|
|
|
brotli_7
|
|
|
brotli_8
|
|
|
brotli_9
|
|
|
brotli_10
|
|
|
brotli_11
|
|
|
zlib_1
|
|
|
zlib_2
|
|
|
zlib_3
|
|
|
zlib_4
|
|
|
zlib_5
|
|
|
zlib_6
|
|
|
zlib_7
|
|
|
zlib_8
|
|
|
zlib_9
|
|
|
zstd_1
|
|
|
zstd_2
|
|
|
zstd_3
|
|
|
zstd_4
|
|
|
zstd_5
|
|
|
zstd_6
|
|
|
zstd_7
|
|
|
zstd_8
|
|
|
zstd_9
|
|
|
zstd_10
|
|
|
zstd_11
|
|
|
zstd_12
|
|
|
zstd_13
|
|
|
zstd_14
|
|
|
zstd_15
|
|
|
zstd_16
|
|
|
zstd_17
|
|
|
zstd_18
|
|
|
zstd_19
|
|
|
zstd_20
|
|
|
zstd_21
|
|
|
lzma_0
|
|
|
lzma_1
|
|
|
lzma_2
|
|
|
lzma_3
|
|
|
lzma_4
|
|
|
lzma_5
|
|
|
lzma_6
|
|
|
lzma_7
|
|
|
lzma_8
|
|
|
lzma_9
|
|
|
bzip2_1
|
|
|
bzip2_2
|
|
|
bzip2_3
|
|
|
bzip2_4
|
|
|
bzip2_5
|
|
|
bzip2_6
|
|
|
bzip2_7
|
|
|
bzip2_8
|
|
|
bzip2_9
|
|
|
zstd_fast_1
|
|
|
zstd_fast_2
|
|
|
zstd_fast_3
|
|
|
zstd_fast_4
|
|
|
zstd_fast_5
|
|
|
zstd_fast_6
|
|
|
zstd_fast_7
|
|
|
zlib6
|
|
|
gzip_normal
|
|
|
zlib9
|
|
|
gzip_best_compression
|
|
|
zstd
|
|
|
brotli3
|
|
|
brotli5
|
|
|
brotli8
|
|
|
quick_lz
|
|
NYT::NConcurrency::EExecutionStackKind
|
Possible values
|
Description
|
|
small
|
|
|
large
|
|
|
huge
|
|
NYT::NConcurrency::TFiberManagerConfig
Source: yt/yt/core/concurrency/config.h
NYT::NConcurrency::TFiberManagerDynamicConfig
Source: yt/yt/core/concurrency/config.h
NYT::NCoreDump::TCoreDumperConfig
Source: yt/yt/library/coredumper/config.h
|
Parameter
|
Description
|
|
path
|
Type: std::string
Required parameter
|
|
pattern
|
Type: std::string
Default value: core.%CORE_DATETIME.%CORE_PID.%CORE_SIG.%CORE_THREAD_NAME-%CORE_REASON
|
NYT::NCrypto::TPemBlobConfig
Source: yt/yt/core/crypto/config.h
|
Parameter
|
Description
|
|
environment_variable
|
Type: std::optional<std::string>
|
|
file_name
|
Type: std::optional<std::string>
|
|
value
|
Type: std::optional<std::string>
|
NYT::NCrypto::TSslContextCommand
Source: yt/yt/core/crypto/config.h
|
Parameter
|
Description
|
|
name
|
Type: std::string
Required parameter
|
|
value
|
Type: std::string
Required parameter
|
NYT::NFlow::EBacktraceEnricherLevel
|
Possible values
|
Description
|
|
enabled_for_all
|
|
|
enabled_for_trivial_errors
|
|
|
enabled_for_not_native_errors
|
|
|
disabled
|
|
NYT::NFlow::EBalanceResource
|
Possible values
|
Description
|
|
cpu
|
|
|
memory
|
|
NYT::NFlow::EBalancerMetricsSource
|
Possible values
|
Description
|
|
job
|
Only the running job's metrics: the 10-minute rate, and until it exists the 30-second and immediate ones. The previous behaviour.
|
|
partition
|
The 10-minute rate of the running job, and until it exists the partition's persisted history (measured by the previous job, survives moves and restarts).
|
NYT::NFlow::EClickHouseCodec
|
Possible values
|
Description
|
|
none
|
No compression.
|
|
lz4
|
LZ4 compression: minimal CPU cost at a moderate compression ratio.
|
|
zstd
|
ZSTD compression: a higher compression ratio at the cost of more CPU.
|
NYT::NFlow::EClickHouseHostSelectionPolicy
|
Possible values
|
Description
|
|
ordered_round_robin
|
Preserves the configured endpoint order; this is the default.
|
|
random_start
|
Chooses one uniformly distributed starting position at client construction, then preserves exhaustive round-robin endpoint iteration.
|
NYT::NFlow::EDistributionOrdering
|
Possible values
|
Description
|
|
strict
|
Output messages (created from input messages with the same key) are distributed in the order they were created. This order survives repartitioning (no race condition between stopping and running partitions that distribute outgoing messages originating from incoming messages with the same key).
|
|
relaxed
|
A race condition between stopping and running partitions is allowed. Slightly reduces message processing latency during repartitioning.
|
NYT::NFlow::EFetchType
|
Possible values
|
Description
|
|
select_rows
|
Supports dynamic (including replicated) tables. Does not support read-ahead; queries go directly to tablet nodes (recommended option).
|
|
table_reader
|
Supports dynamic and static tables. Supports read-ahead; queries go through the master node.
|
NYT::NFlow::EFlowStateTarget
|
Possible values
|
Description
|
|
all
|
All state categories (default).
|
|
key_state
|
Only key_states.
|
|
partition_state
|
Only partition_states.
|
|
external_key_state
|
Only external states (external_key_states and, for read-states, joined_external_key_states).
|
NYT::NFlow::EJobBalancerType
|
Possible values
|
Description
|
|
greedy
|
|
|
cpu_aware
|
|
|
resource_queue
|
|
NYT::NFlow::EPipelineState
|
Possible values
|
Description
|
|
unknown
|
Default state value for a pipeline that has not been started yet.
|
|
stopped
|
Pipeline is stopped. Draining has been performed — all intermediate messages in the pipeline are processed, the actual processed offsets in the source queues are committed. All jobs are stopped.
In this state, you can safely roll out a pipeline release and update its static spec.
To initiate a transition to this state, use the stop-pipeline command.
|
|
paused
|
Pipeline is paused. All jobs are stopped, but intermediate messages in the pipeline may not have been processed.
To initiate a transition to this state, use the pause-pipeline command.
|
|
working
|
Pipeline is working. Messages are being processed.
To initiate a transition to this state, use the start-pipeline command.
|
|
draining
|
Intermediate state. Pipeline is in the process of stopping (transition to stopped state). All messages that the pipeline has already seen from sources and all internal intermediate messages between computations are being final-processed.
|
|
pausing
|
Intermediate state. Pipeline is in the process of pausing. All jobs are being stopped.
|
|
completed
|
Final state. Pipeline is complete. All sources were finite, and all messages from them have been processed.
It is currently not possible to exit this state; you can only recreate the pipeline.
|
NYT::NFlow::EProcessingMode
|
Possible values
|
Description
|
|
exactly_once
|
Default value. The result of Transform, including output messages, processed messages, etc., is committed to YTsaurus within a single transaction. All input messages are deduplicated by message_id.
|
|
at_least_once_consistent
|
In this mode, input message deduplication is disabled. In this mode, TransformComputation stops interacting with the input_messages table. Input messages may be processed multiple times, and the processing result is saved to output_messages each time and is guaranteed to be processed by subsequent Computation instances. This mode can be used if duplicates are not a problem, for example, because the user logic itself can deduplicate redundant messages. Duplicates can occur during any job restarts/crashes (including rescheduling). System shutdown via draining using stop-pipeline, however, does not lead to a violation of guarantees (unless jobs are restarted for some other reason during the process).
|
NYT::NFlow::EQueueTabletIndexRoutingHashPolicy
|
Possible values
|
Description
|
|
range
|
The uint64 hash from tablet_index_routing_hash_expression is reduced to contiguous equal-width hash ranges (rangeSize = 2^64 / tablet_count); recommended: a consumer range-partitioned by the same key reads only its own tablet.
|
|
modulo
|
The uint64 hash from tablet_index_routing_hash_expression is reduced as hash % tablet_count. Discouraged: Flow computations are range-partitioned by key, so a modulo-partitioned queue forces every reader to read every tablet (a full mesh on read). Use range instead.
|
NYT::NFlow::ETimeType
|
Possible values
|
Description
|
|
event_time
|
|
|
system_time
|
|
|
current_time
|
|
|
Possible values
|
Description
|
|
seconds
|
|
|
milli_seconds
|
|
|
iso8601
|
|
NYT::NFlow::EUnavailableSourcePolicy
|
Possible values
|
Description
|
|
retry
|
Default value. A source read error drops the traversal iteration; the read is retried. The traversal does not advance past the unread range.
|
|
mark_unreadable
|
A source read error is swallowed. The traversal continues, and the keys of the unread range are resolved with an uninitialized accessor (IsInitialized() == false).
|
NYT::NFlow::EWorkerCoefMode
|
Possible values
|
Description
|
|
legacy
|
From the current usage of each worker's partitions against the computation averages. The previous behaviour.
|
|
probing
|
From partitions that moved between workers: CPU per message before and after the move.
|
NYT::NFlow::NCompanion::TCompanionConfig
Source: yt/yt/flow/library/cpp/companion/client/config.h
|
Parameter
|
Description
|
|
port
|
Type: int
Default value: 0
The port on which the worker talks to the companion process. The companion does not start without it. For a vanilla launch it is filled in automatically: 10082 while the worker task runs on fixed ports (port_count = 0), or the value of YT_PORT_2 with port_count = 3.
|
|
monitoring_port
|
Type: int
Default value: 0
|
|
companion_process_count
|
Type: int
Default value: 0
|
|
http_client_config
|
Type: NYT::TIntrusivePtr<NYT::NHttp::TClientConfig>
Default value: {}
Config of the HTTP client a C++ companion hands to process functions through IRuntimeInitContext::GetHttpClient(). Mirrors the same-named TFlowNodeConfig field; companions in other languages ignore it.
|
|
https_client_config
|
Type: NYT::TIntrusivePtr<NYT::NHttps::TClientConfig>
Default value: {}
Config of the HTTPS client a C++ companion hands to process functions through IRuntimeInitContext::GetHttpsClient(). Mirrors the same-named TFlowNodeConfig field; companions in other languages ignore it. The config reaches the companion through the process environment. For credentials.private_key, only file_name is supported; inline and environment-backed private keys are rejected.
|
|
http_poller_threads
|
Type: int
Default value: 1
Thread count of the C++ companion's HTTP poller, which runs its HTTP and HTTPS clients. Companions in other languages ignore it.
|
NYT::NFlow::NController::TControllerConfig
Source: yt/yt/flow/library/cpp/controller/config.h
|
Parameter
|
Description
|
|
controller_threads
|
Type: int
Default value: 5
Number of threads.
|
|
orchid_update_period
|
Type: TDuration
Default value: 1s
orchid recalculation period.
|
|
warm_up_time
|
Type: TDuration
Default value: 5s
Warm-up period when the controller starts.
|
|
scheduler_period
|
Type: TDuration
Default value: 5s
Scheduler invocation period.
|
|
cache_period
|
Type: TDuration
Default value: 1s
|
|
feedback_period
|
Type: TDuration
Default value: 1s
|
|
metrics_period
|
Type: TDuration
Default value: 5s
|
|
write_own_retryable_errors_period
|
Type: TDuration
Default value: 5s
|
|
publish_retry_period
|
Type: TDuration
Default value: 5s
Period during which the controller attempts to tell YTsaurus that it is the leader.
|
|
publish_timeout
|
Type: TDuration
Default value: 2h
|
|
election_manager
|
Type: NYT::NYTree::TPolymorphicYsonStruct<NYT::NYTree::NDetail::TPolymorphicMapping<&NYT::NFlow::NController::ElectionBackendDiscriminator.<char const at offset 0>, NYT::NFlow::NController::EElectionBackend, NYT::NYTree::NDetail::TOptionalValue<NYT::NFlow::NController::EElectionBackend, (NYT::NFlow::NController::EElectionBackend)0>, NYT::NFlow::NController::TElectionBackendConfigBase, NYT::NYTree::NDetail::TLeafTag<(NYT::NFlow::NController::EElectionBackend)0, NYT::NFlow::NController::TCypressElectionBackendConfig>, NYT::NYTree::NDetail::TLeafTag<(NYT::NFlow::NController::EElectionBackend)1, NYT::NFlow::NController::TDyntableElectionBackendConfig>, NYT::NYTree::NDetail::TLeafTag<(NYT::NFlow::NController::EElectionBackend)2, NYT::NFlow::NController::TChaosElectionBackendConfig>>>
Default value:
{
"backend" = "cypress";
"leader_cache_update_period" = 1000;
"leader_lease_ping_period" = 1000;
"leader_lease_ttl" = 5000;
"lock_acquisition_period" = 1000;
}
Leader election settings for multiple controllers.
|
|
persisted_state_manager
|
Type: NYT::TIntrusivePtr<NYT::NFlow::NController::TPersistedStateManagerConfig>
Default value: {}
PersistedStateManager settings.
|
|
lease_manager
|
Type: NYT::TIntrusivePtr<NYT::NFlow::NController::TLeaseManagerConfig>
Default value: {}
LeaseManager settings.
|
|
controller_service
|
Type: NYT::TIntrusivePtr<NYT::NFlow::NController::TControllerServiceConfig>
Default value: {}
ControllerService settings.
|
|
bus
|
Type: NYT::TIntrusivePtr<NYT::NBus::NTcp::TBusConfig>
Default value: {}
|
Additional parameters
|
publish_request_timeout
|
Type: TDuration
Default value: 1s
Timeout of a single request of a publication attempt. A failed attempt is retried, so this only decides how long one attempt may hang before the retry.
|
NYT::NFlow::NController::TControllerServiceConfig
Source: yt/yt/flow/library/cpp/controller/config.h
|
Parameter
|
Description
|
|
set_spec_retry_count
|
Type: int
Default value: 3
Number of attempts to update the spec in case of errors.
|
|
set_spec_retry_period
|
Type: TDuration
Default value: 5s
Time between attempts.
|
|
tables_throttler
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TLoadThroughputThrottlerSpec>
Default value: {}
|
NYT::NFlow::NController::TLeaseManagerConfig
Source: yt/yt/flow/library/cpp/controller/config.h
|
Parameter
|
Description
|
|
lease_timeout
|
Type: TDuration
Default value: 10m
How long a job stays fenced without the controller refreshing the fence. With the Cypress election backend this is the timeout of the job's master lease transaction, with the chaos one the timeout of its chaos lease, and with the dyntable one the ttl of the pipeline-wide deadline row that gates every worker commit.
|
|
lease_ping_period
|
Type: TDuration
Default value: 30s
Period at which the leader prolongs the leases. Every backend prolongs from the leader and none from the worker, but the machinery differs: the chaos leases are pinged by one periodic executor, the Cypress ones by a client-side pinger per lease transaction, and the shared dyntable deadline is simply rewritten.
|
|
max_concurrent_requests
|
Type: long
Default value: 500
Maximum number of concurrent requests the manager keeps in flight when it attaches to leases or terminates them. Do not increase to avoid overloading the master. Prolongation is not capped by it: the chaos pings are dispatched in one fire-and-forget round, and the Cypress ones belong to the lease transactions themselves.
|
NYT::NFlow::NController::TPersistedStateManagerConfig
Source: yt/yt/flow/library/cpp/controller/config.h
|
Parameter
|
Description
|
|
timeout
|
Type: TDuration
Default value: 5s
Save timeout in YTsaurus.
|
|
max_reads_per_transaction
|
Type: long
Default value: 10000
|
|
max_writes_per_transaction
|
Type: long
Default value: 10000
|
NYT::NFlow::NDeltaCodecs::ECodec
|
Possible values
|
Description
|
|
none
|
|
|
x_delta
|
|
|
v_c_diff
|
|
NYT::NFlow::NFileStorage::TFileStorageConfig
Source: yt/yt/flow/library/cpp/file_storage/config.h
|
Parameter
|
Description
|
|
path
|
Type: std::string
Required parameter
Exact absolute cache root for one worker process. Flow does not append a pipeline, operation, or job-cookie identifier. The worker holds an exclusive nonblocking lock on <path>/.lock and fails at startup if the same physical root is already in use.
|
|
soft_size_limit
|
Type: long
Required parameter
Target total size of managed payload files and not-yet-deleted trash in bytes after LRU cleanup of unpinned objects.
|
|
hard_size_limit
|
Type: long
Required parameter
Admission limit in bytes for managed payload files, not-yet-deleted trash, and expected-size reservations, including pinned objects. It is not a physical volume quota and cannot strictly cover unknown or misreported staging sizes and filesystem overhead.
|
|
cleanup_period
|
Type: TDuration
Default value: 1m
Period between background cache reconciliation and cleanup of unpinned objects toward soft_size_limit.
|
NYT::NFlow::NStaticTableConnector::TTableTimestampLocatorSpec
Source: yt/yt/flow/library/cpp/connectors/static_table/source_spec.h
|
Parameter
|
Description
|
|
attribute
|
Type: std::string
Required parameter
|
|
format
|
Type: NYT::NFlow::ETimestampFormat
Default value: iso8601
|
|
timezone
|
Type: std::optional<std::string>
Optional IANA time zone used to interpret ISO 8601 timestamps without an explicit UTC offset. Explicit offsets take precedence. The parameter can only be used with format=iso8601. If omitted, timestamps without an offset are interpreted as UTC for backward compatibility.
|
NYT::NFlow::NWorker::TWorkerConfig
Source: yt/yt/flow/library/cpp/worker/config.h
NYT::NFlow::TAtMostOnceStrategyDynamicParameters
Source: yt/yt/flow/library/cpp/connectors/common/async_at_most_once_sink_base.h
|
Parameter
|
Description
|
|
suspend_destruction_duration
|
Type: TDuration
Default value: 10s
|
|
total_queue_bytes_limit
|
Type: NYT::NYTree::TSize
Default value: 100Mi
|
NYT::NFlow::TAtMostOnceStrategyParameters
Source: yt/yt/flow/library/cpp/connectors/common/async_at_most_once_sink_base.h
|
Parameter
|
Description
|
|
enabled
|
Type: bool
Default value: false
|
NYT::NFlow::TAuthenticatorConfig
Source: yt/yt/flow/library/cpp/common/authenticator.h
|
Parameter
|
Description
|
|
require_proxy_signature
|
Type: bool
Default value: false
|
NYT::NFlow::TBacktraceEnricherDynamicSpec
Source: yt/yt/flow/library/cpp/misc/error_backtrace_enricher.h
NYT::NFlow::TBacktraceEnricherSpec
Source: yt/yt/flow/library/cpp/misc/error_backtrace_enricher.h
NYT::NFlow::TComputationSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
computation_class_name
|
Type: std::string
Required parameter
The built-in Computation or process-function adapter class name.
For new C++ user logic, select one of the built-in TProcessFunction*Computation adapters. Register the function with YT_FLOW_DEFINE_PROCESS_FUNCTION and name it in processing_function. Don’t create custom Computation classes.
|
|
processing_function
|
Type: std::optional<std::string>
The fully qualified process-function name registered with YT_FLOW_DEFINE_PROCESS_FUNCTION. Required by TProcessFunction*Computation adapters.
|
|
processing_function_parameters
|
Type: NYT::TIntrusivePtr<NYT::NYTree::IMapNode>
Static process-function parameters. Their type is declared in YT_FLOW_DEFINE_PROCESS_FUNCTION registration.
|
|
group_by_schema
|
Type: NYT::TIntrusivePtr<NYT::NTableClient::TTableSchema>
Default value: {'value': [], 'attributes': {'strict': true, 'unique_keys': false}}
A schema for grouping all input streams. It is a dynamic table schema. Schema properties and requirements:
- Does not contain any ordering rules.
- Can contain computed fields.
- Must be empty only for
SourceComputation.
- For other
Computation instances, the first column must be of type uint64. It is recommended that this column be a computed hash of the other columns (for example, farm_hash). Note that using the modulo operator % and bigb_hash in expression is prohibited.
|
|
experimental_enable_non_uint_key
|
Type: std::optional<bool>
|
|
input_stream_ids
|
Type: THashSet<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TStreamIdTag>>
Default value: []
Input streams.
The schemas of all streams must be castable to group_by_schema.
|
|
output_stream_ids
|
Type: THashSet<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TStreamIdTag>>
Default value: []
Output streams of the Computation. They can be used as inputs for other Computations.
|
|
streams_dependency
|
Type: THashMap<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TStreamIdTag>, THashSet<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TStreamIdTag>>>
Default value: {}
A description of dependencies between streams, including internal source and timer streams.
Example: {stream1, {stream2, stream3}} — stream1 depends on streams stream2 and stream3.
Specified only for timer and output streams. If no description is provided for a stream, the dependency list is auto-generated:
- for
timer streams, the list includes all input and source streams;
- for
output streams, the list includes all timer, input, and source streams.
|
|
watermark_strategy
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TWatermarkStrategySpec>
Default value: {}
Settings related to EventWatermark.
An optional parameter. Must be specified within SourceComputation.
|
|
required_resource_ids
|
Type: THashMap<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TResourceIdTag>, NYT::TIntrusivePtr<NYT::NFlow::TResourceDescription>>
Default value: {}
A list of resources required for the Computation to work.
Has an alias parameter for locally renaming a resource (for example, Computation expects a resource named YTClient, but the provided resource has a different name). The controller and worker parameters control whether the resource is created on the Controller or on the Worker.
|
|
parameters
|
Type: NYT::TIntrusivePtr<NYT::NYTree::IMapNode>
Default value: {}
Custom parameters for Computation and ComputationController.
|
|
timer_streams
|
Type: THashMap<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TStreamIdTag>, NYT::TIntrusivePtr<NYT::NFlow::TTimerSpec>>
Default value: {}
Settings for all timer streams. Keys must match [0-9A-Za-z_-]+.
|
|
key_visitor_streams
|
Type: THashMap<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TStreamIdTag>, NYT::TIntrusivePtr<NYT::NFlow::TKeyVisitorStreamSpec>>
Default value: {}
Settings for all key_visitor streams. Keys must match [0-9A-Za-z_-]+.
|
|
source_streams
|
Type: THashMap<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TStreamIdTag>, NYT::TIntrusivePtr<NYT::NFlow::TSourceSpec>>
Default value: {}
Settings for all source streams. Keys must match [0-9A-Za-z_-]+.
|
|
sinks
|
Type: THashMap<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TSinkIdTag>, NYT::TIntrusivePtr<NYT::NFlow::TSinkSpec>>
Default value: {}
Settings for all sinks. Keys must match [0-9A-Za-z_-]+.
|
|
external_state_managers
|
Type: THashMap<std::string, NYT::TIntrusivePtr<NYT::NFlow::TExternalStateManagerSpec>>
Default value: {}
A declarative declaration of external state managers for this Computation. The key is the client name that the process function passes to InitExternalStateClient (must start with /, for example /state); the value contains the manager class name and its parameters.
|
|
external_state_joiners
|
Type: THashMap<std::string, NYT::TIntrusivePtr<NYT::NFlow::TExternalStateJoinerSpec>>
Default value: {}
A declarative declaration of external state joiners (read-only access to external states via key join) for this Computation. The key is the client name that the process function passes to InitExternalStateClient (must start with /, for example /state); the value contains the joiner class name and its parameters.
|
|
state_joiners
|
Type: THashMap<std::string, NYT::TIntrusivePtr<NYT::NFlow::TStateJoinerSpec>>
Default value: {}
A declarative declaration of state joiners (read-only access to the internal state of another Computation via key join) for this Computation. The key is the client name that the process function passes to InitClient (must start with /); the value specifies the target computation_id, its state_name, and join_on.
|
|
heavy_hitters
|
Type: NYT::NFlow::THeavyHittersSpec
Default value: {}
Settings for detecting high-frequency keys.
|
|
pivot_finder
|
Type: NYT::NFlow::TPivotFinderSpec
Default value: {}
|
|
input_ordering
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TInputOrderingSpec>
Default value: {}
Settings for the processing order of input streams.
|
|
distribution_ordering
|
Type: NYT::NFlow::EDistributionOrdering
Default value: strict
Settings for the distribution order of output messages. Controls whether a race condition is possible between stopping and running partitions from the perspective of the order of output messages generated from input messages with the same key.
|
|
worker_group
|
Type: NYT::TStrongTypedef<std::string, NYT::NFlow::TWorkerGroupIdTag, NYT::TStrongTypedefOptions{true}>
Default value: TString("")
|
|
allow_timer_self_dependency
|
Type: bool
Default value: false
|
|
use_compact_input_messages
|
Type: std::optional<bool>
Use the compact_input_messages table for input message deduplication instead of input_messages. The compact deduplication key stores only the hash of key[0] and CityHash128(message_id), which reduces the row size (~200 bytes → ~71 bytes).
An optional parameter. If not set, the table is chosen automatically: the compact table is always used except when experimental_enable_non_uint_key is enabled (in that case, the first key column may not be uint64, and deduplication requires the full key). An explicitly set value takes priority over automatic selection. Enabling use_compact_input_messages together with experimental_enable_non_uint_key is prohibited.
|
NYT::NFlow::TDeleteStatesArg
Source: yt/yt/flow/library/cpp/controller/state_access.h
Argument of the delete-states command.
The same three modes as in TReadStatesArg are supported.
By default, it is a dry run: the request returns counters of matched rows without deleting anything. For actual deletion, you must pass commit=true.
The command requires the pipeline to be in Stopped/Completed state, or force=true when in Paused. Only key/partition/manager states are deleted; joiner states are left untouched.
|
Parameter
|
Description
|
|
computation_id
|
Type: std::optional<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TComputationIdTag>>
ID of the Computation. Required for modes 1 and 2 (see the command description).
|
|
partition_id
|
Type: std::optional<NYT::TStrongTypedef<NYT::TGuid, NYT::NFlow::TPartitionIdTag, NYT::TStrongTypedefOptions{true}>>
Partition ID. Required for mode 3 (see the command description).
|
|
key
|
Type: std::optional<NYT::NYson::TYsonString>
Key for point lookup (mode 2). Accepted in two forms:
- YSON Dict
{column = value; ...} — for each target table (key_states of the Computation and each table of the external_state_manager/joiner with its own key_schema_override) the dict is decomposed by its key schema; columns absent from the schema are ignored.
- YSON List of positional values — passed to tables as-is, without adaptation to their schema.
|
|
name
|
Type: std::optional<std::string>
Filter by exact state name; applies to all sections (key/partition/external). If set, only rows with this state name are returned in each section.
|
|
target
|
Type: NYT::NFlow::EFlowStateTarget
Default value: all
Which categories of states to request. Default: all.
|
|
force
|
Type: bool
Default value: false
Allows deletion on a pipeline in Paused state. Without force, the command requires Stopped/Completed.
|
|
commit
|
Type: bool
Default value: false
If true, the found rows are actually deleted. Otherwise, the request works in dry-run mode and only counts them.
|
NYT::NFlow::TDeleteStatesResponse
Source: yt/yt/flow/library/cpp/controller/state_access.h
Response of the delete-states command.
|
Parameter
|
Description
|
|
committed
|
Type: bool
Required parameter
Reflects the commit value from the request argument. If errors contains entries, some matched rows might not have been deleted — exact counters are not provided; in practice, the state must be re-read via read-states.
|
|
matched_states
|
Type: NYT::NFlow::TMatchedStates
Required parameter
Counters of matched rows by category. Always populated: both in dry-run mode (shows what would have been deleted) and with commit=true (shows what was deleted or attempted to be deleted). See TMatchedStates.
|
|
errors
|
Type: std::vector<std::string>
Default value: []
Errors that occurred while forming the response to the request itself (for example, a Sync manager failure). These are not related to errors in the pipeline's work with the state. If such an error occurred for one of the external states, other categories can still be processed successfully.
|
NYT::NFlow::TDirectControllerCommandsConfig
Source: yt/yt/flow/library/cpp/pipeline_helpers/flow_execute/flow_execute.h
|
Parameter
|
Description
|
|
enabled
|
Type: bool
Default value: false
|
|
rpc_timeout
|
Type: TDuration
Default value: 30s
|
NYT::NFlow::TDynamicBufferStateManagerSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
manage_period
|
Type: TDuration
Default value: 1s
Frequency of buffer size recalculation.
|
|
demand_window
|
Type: TDuration
Default value: 1m
Time window over which the utilization of a single buffer is estimated.
|
|
input_buffer
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TDynamicBufferStateManagerSpec::TOneSideBufferSpec>
Default value: {}
Input buffer settings.
|
|
output_buffer
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TDynamicBufferStateManagerSpec::TOneSideBufferSpec>
Default value: {}
Output buffer settings.
|
|
enable_v2
|
Type: bool
Default value: true
Enables the v2 buffer-sizing strategy: limit = peak usage + headroom, with a v2_gain_epochs × demand × epoch bandwidth-delay-product floor. The demand × max_duration cap is raised by announced input backlog or, for a producing output, to an equal share of half fair_share_pool because producer epochs do not reveal the downstream acknowledgement period. Demand-backed limits are allocated before speculative output probes. On both sides, Σ(limits) ≤ fair_share_pool including other streams' in-flight bytes; job_limit remains the per-stream bound. Streams with job_overrides remain fully outside the pool as in v1, so worker buffer memory should cover fair_share_pool plus the actual in-flight bytes of overridden streams. Enabled by default; set this option to false to use the previous v1 formula.
|
Additional parameters
|
epoch_cycle_window_samples
|
Type: int
Default value: 16
Number of samples in the window used to estimate the median job epoch cycle.
|
|
max_rate_estimator_buckets
|
Type: int
Default value: 8
Number of buckets in the windowed-max drain-rate estimator; changing it resets the estimate.
|
|
warmup_refresh_period
|
Type: TDuration
Default value: 30s
How often a job polls the converged warmup state for persistence.
|
|
v2_gain_epochs
|
Type: double
Default value: 2.0
Minimum buffer target in units of job epochs, used as a bandwidth-delay-product floor. Applies only when enable_v2 is set.
|
|
v2_use_offered_rate
|
Type: bool
Default value: true
Whether announced backlog rate contributes to demand. Disable it when a producer systematically overstates its offered rate. Applies only when enable_v2 is set.
|
|
v2_floor
|
Type: NYT::NYTree::TSize
Default value: 2Mi
Minimum grant for a stream with backlog; it must fit the largest message. Applies only when enable_v2 is set.
|
|
v2_headroom_growth_factor
|
Type: double
Default value: 2.0
Headroom growth factor per management tick while utilization is high. Applies only when enable_v2 is set.
|
|
v2_high_utilization_threshold
|
Type: double
Default value: 0.65
Utilization threshold above which headroom grows; below half this value headroom decays. Applies only when enable_v2 is set.
|
|
v2_publish_threshold
|
Type: double
Default value: 0.25
Suppress a new limit when its relative increase is smaller than this value; decreases are always published. Applies only when enable_v2 is set.
|
NYT::NFlow::TDynamicBufferStateManagerSpec::TOneSideBufferSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
fair_share_pool
|
Type: NYT::NYTree::TSize
Default value: 3Gi
Pool size for distribution using the FairShare algorithm based on utilization.
|
|
worker_group_fair_share_pool_overrides
|
Type: THashMap<NYT::TStrongTypedef<std::string, NYT::NFlow::TWorkerGroupIdTag, NYT::TStrongTypedefOptions{true}>, NYT::NYTree::TSize>
Default value: {}
Replaces fair_share_pool for workers in the listed worker groups, for installations where some workers have substantially more memory. A worker in several listed groups uses the maximum value, so the pool reflects the memory actually available to that worker.
|
|
job_guarantee
|
Type: NYT::NYTree::TSize
Default value: 5Mi
Minimum buffer size for a single job.
|
|
job_limit
|
Type: NYT::NYTree::TSize
Default value: 500Mi
Maximum buffer size for a single job.
|
|
max_duration
|
Type: TDuration
Default value: 1m
Assume that a buffer need not hold more than max_duration of job work (the rate is estimated heuristically). With enable_v2, it also caps the measured job epoch-cycle estimate. A producing output may probe up to its equal share of half fair_share_pool because producer epochs do not reveal the downstream acknowledgement period. job_limit and the worker pool remain hard bounds.
|
|
job_overrides
|
Type: THashMap<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TComputationIdTag>, THashMap<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TStreamIdTag>, NYT::NYTree::TSize>>
Default value: {}
Ability to manually override the buffer size for a computation stream.
|
NYT::NFlow::TDynamicComputationSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
batch_duration
|
Type: TDuration
Default value: 1s
Maximum time to collect an input batch.
|
|
max_rows_per_batch
|
Type: NYT::NYTree::TSize
Default value: 1K
Maximum size of an input batch in rows. Calculated separately for input_streams, timer_streams, and for each of the source_streams.
|
|
max_bytes_per_batch
|
Type: NYT::NYTree::TSize
Default value: 10Mi
Maximum size of an input batch in bytes. Calculated similarly to max_rows_per_batch.
|
|
max_keys_per_batch
|
Type: std::optional<NYT::NYTree::TSize>
|
|
draining
|
Type: bool
Default value: false
The Draining mode is used to drain the pipeline. In this mode, the Computation stops running timers and fetching new events from the Source.
The stop-pipeline command stops the entire pipeline through a full drain. You can set this mode for individual Computation nodes for debugging.
|
|
empty_batch_backoff
|
Type: TDuration
Default value: 250ms
If the epoch is empty, the Computation sleeps for the specified additional time. We do not recommend setting this parameter to less than 100 ms.
|
|
lease_check_period
|
Type: TDuration
Default value: 1m
|
|
retryable_request
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TDynamicRetryableRequestSpec>
Default value: {}
|
|
tracer
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TDynamicPartitionTracerSpec>
Default value: {}
|
|
parameters
|
Type: NYT::TIntrusivePtr<NYT::NYTree::IMapNode>
Default value: {}
Arbitrary dynamic parameters of the Computation class.
|
|
processing_function_parameters
|
Type: NYT::TIntrusivePtr<NYT::NYTree::IMapNode>
Dynamic process-function parameters. Their type is declared in YT_FLOW_DEFINE_PROCESS_FUNCTION registration.
|
|
source_streams
|
Type: THashMap<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TStreamIdTag>, NYT::TIntrusivePtr<NYT::NFlow::TDynamicSourceSpec>>
Default value: {}
Dynamic parameters of all Sources.
|
|
key_visitor_streams
|
Type: THashMap<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TStreamIdTag>, NYT::TIntrusivePtr<NYT::NFlow::TDynamicKeyVisitorStreamSpec>>
Default value: {}
|
|
sinks
|
Type: THashMap<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TSinkIdTag>, NYT::TIntrusivePtr<NYT::NFlow::TDynamicSinkSpec>>
Default value: {}
Dynamic parameters of all Sinks.
|
|
external_state_managers
|
Type: THashMap<std::string, NYT::TIntrusivePtr<NYT::NFlow::TDynamicExternalStateManagerSpec>>
Default value: {}
|
|
external_state_joiners
|
Type: THashMap<std::string, NYT::TIntrusivePtr<NYT::NFlow::TDynamicExternalStateJoinerSpec>>
Default value: {}
|
|
state_joiners
|
Type: THashMap<std::string, NYT::TIntrusivePtr<NYT::NFlow::TDynamicStateJoinerSpec>>
Default value: {}
|
|
state_manager
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TDynamicStateManagerSpec>
Default value: {}
|
|
input_store
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TDynamicInputStoreSpec>
Default value: {}
|
|
timer_store
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TDynamicTimerStoreSpec>
Default value: {}
|
|
output_store
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TDynamicOutputStoreSpec>
Default value: {}
|
|
timer_store_count_limit
|
Type: NYT::NYTree::TSize
Default value: 250K
|
|
timer_store_byte_size_limit
|
Type: NYT::NYTree::TSize
Default value: 100Mi
|
|
output_store_count_limit
|
Type: NYT::NYTree::TSize
Default value: 1G
|
|
output_store_byte_size_limit
|
Type: NYT::NYTree::TSize
Default value: 100Gi
|
|
blocked_time_window
|
Type: TDuration
Default value: 10m
Averaging window for the share of time the job spent blocked on each limit (blocked_share in diagnostics). The share is normalized to the job lifetime, so a job blocked from the start shows ~1 regardless of its age.
|
|
input_rows_throttler_id
|
Type: std::optional<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TThrottlerIdTag>>
Throttler ID for limiting the message processing rate. If set, before each iteration the Computation waits for a quota equal to the number of messages in the input batch. The ID must be present in dynamic_spec/throttlers. For more information, see Distributed Throttler.
|
|
input_bytes_throttler_id
|
Type: std::optional<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TThrottlerIdTag>>
Similar to input_rows_throttler_id, but quotas the total byte_size of messages in the batch — the system size of their serialized representation.
|
|
input_rows_throttler_class_id
|
Type: std::optional<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TQuotaClassIdTag>>
Quota class for the input_rows_throttler_id throttler. The class must be declared by that throttler; default selects the reserved class with weight 1.0. Requires input_rows_throttler_id to be set.
|
|
input_bytes_throttler_class_id
|
Type: std::optional<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TQuotaClassIdTag>>
Quota class for the input_bytes_throttler_id throttler. The class must be declared by that throttler; default selects the reserved class with weight 1.0. Requires input_bytes_throttler_id to be set.
|
|
skip_if_expression
|
Type: std::optional<std::string>
YTQL predicate for filtering (skipping) input messages before Computation processing. A message is skipped (dropped) if the predicate evaluates to boolean true.
The predicate is evaluated on the message payload columns plus the meta-columns $message_id, $stream_id, $system_timestamp, $event_timestamp, $alignment_timestamp. The result must be boolean: NULL or a non-boolean result causes an exception.
Applies to both input and source streams. The number of skipped messages is exported via the input_streams/skipped_by_expression_count and source_streams/skipped_by_expression_count metrics.
|
NYT::NFlow::TDynamicControllerConnectorSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
controller_wait_timeout
|
Type: TDuration
Default value: 5m
Time after which the worker forcibly disconnects all jobs in case of loss of connection with the controller.
|
|
controller_discover_period
|
Type: TDuration
Default value: 10s
Period for rediscovering the controller.
|
|
controller_heartbeat_period
|
Type: TDuration
Default value: 1s
Period of worker heartbeats to the controller.
|
|
worker_statistics_report_period
|
Type: TDuration
Default value: 30s
|
|
controller_heartbeat_rpc_timeout
|
Type: TDuration
Default value: 10s
Heartbeat timeout.
|
|
controller_heartbeat_failure_backoff
|
Type: TDuration
Default value: 1s
Backoff time in case of handshake failure.
|
|
controller_handshake_rpc_timeout
|
Type: TDuration
Default value: 10s
Handshake timeout.
|
|
controller_handshake_failure_backoff
|
Type: TDuration
Default value: 1s
|
|
orchid_update_period
|
Type: TDuration
Default value: 1s
Period for updating orchid on the controller.
|
NYT::NFlow::TDynamicExpiringJobNamedStateCacheSpec
Source: yt/yt/flow/library/cpp/common/state_cache.h
|
Parameter
|
Description
|
|
ttl
|
Type: TDuration
Default value: 0
|
NYT::NFlow::TDynamicExternalStateJoinerSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
parameters
|
Type: NYT::TIntrusivePtr<NYT::NYTree::IMapNode>
Dynamic parameters of the specific joiner implementation.
|
NYT::NFlow::TDynamicExternalStateManagerSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
parameters
|
Type: NYT::TIntrusivePtr<NYT::NYTree::IMapNode>
Dynamic parameters of the specific manager implementation.
|
NYT::NFlow::TDynamicFileProviderSpec
Source: yt/yt/flow/library/cpp/common/file_provider.h
|
Parameter
|
Description
|
|
parameters
|
Type: NYT::TIntrusivePtr<NYT::NYTree::IMapNode>
Default value: {}
Dynamic parameters of the registered file-provider implementation. They select future revisions during discovery but do not change an exact revision that has already been delivered.
|
Source: yt/yt/flow/library/cpp/common/spec.h
This structure has no main parameters.
NYT::NFlow::TDynamicJobManagerGroupSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
balancer_type
|
Type: NYT::NFlow::EJobBalancerType
Default value: cpu_aware
|
|
use_cpu_aware_balancer
|
Type: std::optional<bool>
Enables experimental CPU-aware balancing. Default balancing can lead to severe CPU resource imbalance among workers.
|
|
rebalance_delay_after_pipeline_sync
|
Type: TDuration
Default value: 30s
|
|
rebalance_target_deviation
|
Type: double
Default value: 0.1
Maximum acceptable deviation.
|
|
rebalance_hot_mode_coeff
|
Type: double
Default value: 2.0
|
|
rebalance_action_min_time
|
Type: TDuration
Default value: 4s
|
|
rebalance_action_max_time
|
Type: TDuration
Default value: 12s
|
|
rebalance_sync_period
|
Type: TDuration
Default value: 10s
Period of cpu_aware balancing. It is not recommended to set a value lower than 15min because the balancer looks at 10min metrics.
|
|
rebalance_count_exceeded_allowed
|
Type: double
Default value: 1.2
|
|
rebalance_even_load_thresholds
|
Type: THashMap<NYT::NFlow::EBalanceResource, NYT::TIntrusivePtr<NYT::NFlow::TEvenLoadThresholds>>
Default value: {}
Per-resource thresholds of the "load is even" gate, for example {cpu = {spread = 2.0; ratio = 1.5}}. Only resources with a non-zero weight are consulted. Unset fields fall back to the defaults: a spread of 1.0 core for CPU and 1 GB for memory, a ratio of 1.2.
|
|
rebalance_min_cpu_spread
|
Type: std::optional<double>
Deprecated: the same as spread of the cpu entry of rebalance_even_load_thresholds. A value set explicitly in rebalance_even_load_thresholds takes precedence.
|
|
rebalance_min_cpu_ratio
|
Type: std::optional<double>
Deprecated: the same as ratio of the cpu entry of rebalance_even_load_thresholds. A value set explicitly in rebalance_even_load_thresholds takes precedence.
|
|
balance_weights
|
Type: THashMap<NYT::NFlow::EBalanceResource, double>
Default value: {'cpu': 1.0, 'memory': 0.0}
Relative importance of the resources for cpu_aware balancing, for example {cpu = 80; memory = 20}. The given keys are merged into the default {cpu = 1; memory = 0}; the weights are normalized, so only the proportions matter. A resource with a zero weight does not take part in balancing, so by default only CPU is balanced. The rebalance_target_deviation threshold applies to the normalized weighted sum: adding a second resource proportionally shrinks the contribution of the first one.
|
|
balancer_metrics_source
|
Type: NYT::NFlow::EBalancerMetricsSource
Default value: job
Where the cpu_aware balancer takes a partition's CPU usage from. partition: the 10-minute rate of the running job, and until it exists the partition's persisted history (measured by the previous job, survives moves and restarts); 30-second and immediate rates are never used. job: the previous behaviour, only the running job's metrics with a fallback to the 30-second and immediate rates. Off by default while the change is rolled out pipeline by pipeline; temporary switch.
|
|
balance_warmup_protection
|
Type: bool
Default value: false
Never move a running job whose metrics are not mature yet: the 10-minute rate exists and one more 10-minute window has passed since it appeared. Such a job is still paying for its previous move and its weight is unknown. The protection applies to every balancing resource: a memory-only configuration waits for the CPU rate to mature as well. A worker over the count limit whose kick candidate is immature is skipped for the round. Jobs younger than two minutes stay movable: they have invested nothing yet, and on a pipeline restart the workers do not register all at once. Partitions without a job are placed as usual. Off by default while the change is rolled out pipeline by pipeline; temporary switch.
|
|
balance_warmup_idle_worker_share
|
Type: double
Default value: 0.2
Share of the group's workers without a job of the group at which warm-up protection is lifted: filling idle workers matters more than the metrics of the jobs that would move, e.g. when workers register minutes apart after a pipeline restart. Temporary switch for rollback.
|
|
worker_coef_mode
|
Type: NYT::NFlow::EWorkerCoefMode
Default value: legacy
How the cpu_aware balancer estimates the relative speed of workers. probing: from partitions that moved between workers, comparing CPU per message before and after the move; observations per worker pair accumulate in balancer_state and are solved together; without moves all coefficients are 1. legacy: the previous estimate from the current usage of each worker's partitions against the computation averages. Off by default while the change is rolled out pipeline by pipeline; temporary switch.
|
|
worker_coef_half_life
|
Type: TDuration
Default value: 1d
Half-life of an accumulated worker-pair observation: when a new observation merges in, the old one loses weight in proportion to the time passed. Without new observations the estimate does not change.
|
|
worker_coef_retention
|
Type: TDuration
Default value: 7d
How long the observations of a worker absent from the group are kept; a worker returning earlier gets its coefficient back.
|
|
worker_coef_prior_weight
|
Type: double
Default value: 0.05
Weight of the prior "the worker's coefficient is 1" in observation weight units (a move of one partition weighs 1 / N, N being the number of partitions on the receiving worker). Pins the group's mean coefficient to 1 and damps single observations.
|
|
worker_coef_max_ratio
|
Type: double
Default value: 4.0
Safety limit: worker coefficients are clamped to [1 / max_ratio, max_ratio].
|
|
disable_even_load_gate
|
Type: std::optional<bool>
|
|
async_balancing
|
Type: bool
Default value: true
|
|
graceful_move
|
Type: bool
Default value: true
|
|
zero_queue_latency
|
Type: TDuration
Default value: 1s
|
|
planning_horizon
|
Type: TDuration
Default value: 10m
|
|
preloading_timeout
|
Type: TDuration
Default value: 30m
How long the resource_queue balancer keeps a model preload that is still loading, counted from the request. Until then the preload is not released and the worker counts as a worker of the computations that need the model; after it the usual rules apply.
|
|
minimum_worker_count
|
Type: unsigned long
Default value: 1
Minimum number of workers to distribute the load across.
Must be filled in to avoid out-of-memory errors and overloading individual workers when the system starts. Recommended value is about 60% of the target number of workers (to survive the loss of a subset of workers, for example, if a data center goes offline).
|
|
lost_job_timeout
|
Type: TDuration
Default value: 10s
Time after which a Job is considered lost if no status update reaches the controller. The same timeout is used to consider a worker lost.
|
|
faulty_address_window
|
Type: TDuration
Default value: 5m
Window within which the exponentially decaying counter of connection losses between the controller and a worker is calculated. If this counter exceeds faulty_address_attempts, the worker is considered faulty and is ignored during load balancing. Solves the problem of workers that constantly lose connection to the controller. Also solves the problem of two workers with different incarnation IDs but the same address (when deployment has two pods that think they share the same FQDN).
|
|
faulty_address_attempts
|
Type: unsigned long
Default value: 5
See the description of the faulty_address_window parameter.
|
NYT::NFlow::TDynamicJobManagerSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
balancer_type
|
Type: NYT::NFlow::EJobBalancerType
Default value: cpu_aware
|
|
use_cpu_aware_balancer
|
Type: std::optional<bool>
Enables experimental CPU-aware balancing. Default balancing can lead to severe CPU resource imbalance among workers.
|
|
rebalance_delay_after_pipeline_sync
|
Type: TDuration
Default value: 30s
|
|
rebalance_target_deviation
|
Type: double
Default value: 0.1
Maximum acceptable deviation.
|
|
rebalance_hot_mode_coeff
|
Type: double
Default value: 2.0
|
|
rebalance_action_min_time
|
Type: TDuration
Default value: 4s
|
|
rebalance_action_max_time
|
Type: TDuration
Default value: 12s
|
|
rebalance_sync_period
|
Type: TDuration
Default value: 10s
Period of cpu_aware balancing. It is not recommended to set a value lower than 15min because the balancer looks at 10min metrics.
|
|
rebalance_count_exceeded_allowed
|
Type: double
Default value: 1.2
|
|
rebalance_even_load_thresholds
|
Type: THashMap<NYT::NFlow::EBalanceResource, NYT::TIntrusivePtr<NYT::NFlow::TEvenLoadThresholds>>
Default value: {}
Per-resource thresholds of the "load is even" gate, for example {cpu = {spread = 2.0; ratio = 1.5}}. Only resources with a non-zero weight are consulted. Unset fields fall back to the defaults: a spread of 1.0 core for CPU and 1 GB for memory, a ratio of 1.2.
|
|
rebalance_min_cpu_spread
|
Type: std::optional<double>
Deprecated: the same as spread of the cpu entry of rebalance_even_load_thresholds. A value set explicitly in rebalance_even_load_thresholds takes precedence.
|
|
rebalance_min_cpu_ratio
|
Type: std::optional<double>
Deprecated: the same as ratio of the cpu entry of rebalance_even_load_thresholds. A value set explicitly in rebalance_even_load_thresholds takes precedence.
|
|
balance_weights
|
Type: THashMap<NYT::NFlow::EBalanceResource, double>
Default value: {'cpu': 1.0, 'memory': 0.0}
Relative importance of the resources for cpu_aware balancing, for example {cpu = 80; memory = 20}. The given keys are merged into the default {cpu = 1; memory = 0}; the weights are normalized, so only the proportions matter. A resource with a zero weight does not take part in balancing, so by default only CPU is balanced. The rebalance_target_deviation threshold applies to the normalized weighted sum: adding a second resource proportionally shrinks the contribution of the first one.
|
|
balancer_metrics_source
|
Type: NYT::NFlow::EBalancerMetricsSource
Default value: job
Where the cpu_aware balancer takes a partition's CPU usage from. partition: the 10-minute rate of the running job, and until it exists the partition's persisted history (measured by the previous job, survives moves and restarts); 30-second and immediate rates are never used. job: the previous behaviour, only the running job's metrics with a fallback to the 30-second and immediate rates. Off by default while the change is rolled out pipeline by pipeline; temporary switch.
|
|
balance_warmup_protection
|
Type: bool
Default value: false
Never move a running job whose metrics are not mature yet: the 10-minute rate exists and one more 10-minute window has passed since it appeared. Such a job is still paying for its previous move and its weight is unknown. The protection applies to every balancing resource: a memory-only configuration waits for the CPU rate to mature as well. A worker over the count limit whose kick candidate is immature is skipped for the round. Jobs younger than two minutes stay movable: they have invested nothing yet, and on a pipeline restart the workers do not register all at once. Partitions without a job are placed as usual. Off by default while the change is rolled out pipeline by pipeline; temporary switch.
|
|
balance_warmup_idle_worker_share
|
Type: double
Default value: 0.2
Share of the group's workers without a job of the group at which warm-up protection is lifted: filling idle workers matters more than the metrics of the jobs that would move, e.g. when workers register minutes apart after a pipeline restart. Temporary switch for rollback.
|
|
worker_coef_mode
|
Type: NYT::NFlow::EWorkerCoefMode
Default value: legacy
How the cpu_aware balancer estimates the relative speed of workers. probing: from partitions that moved between workers, comparing CPU per message before and after the move; observations per worker pair accumulate in balancer_state and are solved together; without moves all coefficients are 1. legacy: the previous estimate from the current usage of each worker's partitions against the computation averages. Off by default while the change is rolled out pipeline by pipeline; temporary switch.
|
|
worker_coef_half_life
|
Type: TDuration
Default value: 1d
Half-life of an accumulated worker-pair observation: when a new observation merges in, the old one loses weight in proportion to the time passed. Without new observations the estimate does not change.
|
|
worker_coef_retention
|
Type: TDuration
Default value: 7d
How long the observations of a worker absent from the group are kept; a worker returning earlier gets its coefficient back.
|
|
worker_coef_prior_weight
|
Type: double
Default value: 0.05
Weight of the prior "the worker's coefficient is 1" in observation weight units (a move of one partition weighs 1 / N, N being the number of partitions on the receiving worker). Pins the group's mean coefficient to 1 and damps single observations.
|
|
worker_coef_max_ratio
|
Type: double
Default value: 4.0
Safety limit: worker coefficients are clamped to [1 / max_ratio, max_ratio].
|
|
disable_even_load_gate
|
Type: std::optional<bool>
|
|
async_balancing
|
Type: bool
Default value: true
|
|
graceful_move
|
Type: bool
Default value: true
|
|
zero_queue_latency
|
Type: TDuration
Default value: 1s
|
|
planning_horizon
|
Type: TDuration
Default value: 10m
|
|
preloading_timeout
|
Type: TDuration
Default value: 30m
How long the resource_queue balancer keeps a model preload that is still loading, counted from the request. Until then the preload is not released and the worker counts as a worker of the computations that need the model; after it the usual rules apply.
|
|
minimum_worker_count
|
Type: unsigned long
Default value: 1
Minimum number of workers to distribute the load across.
Must be filled in to avoid out-of-memory errors and overloading individual workers when the system starts. Recommended value is about 60% of the target number of workers (to survive the loss of a subset of workers, for example, if a data center goes offline).
|
|
lost_job_timeout
|
Type: TDuration
Default value: 10s
Time after which a Job is considered lost if no status update reaches the controller. The same timeout is used to consider a worker lost.
|
|
faulty_address_window
|
Type: TDuration
Default value: 5m
Window within which the exponentially decaying counter of connection losses between the controller and a worker is calculated. If this counter exceeds faulty_address_attempts, the worker is considered faulty and is ignored during load balancing. Solves the problem of workers that constantly lose connection to the controller. Also solves the problem of two workers with different incarnation IDs but the same address (when deployment has two pods that think they share the same FQDN).
|
|
faulty_address_attempts
|
Type: unsigned long
Default value: 5
See the description of the faulty_address_window parameter.
|
|
worker_group_override
|
Type: THashMap<NYT::TStrongTypedef<std::string, NYT::NFlow::TWorkerGroupIdTag, NYT::TStrongTypedefOptions{true}>, NYT::TIntrusivePtr<NYT::NFlow::TDynamicJobManagerGroupSpec>>
Default value: {}
|
|
partition_history_limit
|
Type: long
Default value: 4096
Most partition histories the balancer keeps across job restarts, see balancer_metrics_source. The histories live in one persisted document with a size limit, so when the limit is reached the lightest history gives way to a heavier one and lighter ones are not saved. A pipeline stop saves no histories at all.
|
NYT::NFlow::TDynamicJobTrackerSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
job_control_threads
|
Type: int
Default value: 3
Number of control threads.
In most cases, the default value is sufficient.
|
|
job_threads
|
Type: std::optional<int>
Number of threads for performing jobs.
If not specified, the pool size is calculated automatically based on the node's CPU limit.
If the node limit cannot be calculated, 30 threads are created.
|
|
buffer_state_manager
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TDynamicBufferStateManagerSpec>
Default value: {}
Configuration of the buffer management module for incoming and outgoing messages.
|
|
load_throughput_throttler
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TLoadThroughputThrottlerSpec>
Default value: {'limit': 134217728.0}
|
|
state_cache
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TDynamicStateCacheSpec>
Default value: {}
|
|
mark_performance_metrics_steady_after_first_iteration
|
Type: bool
Default value: false
Report the job's rate counters (CPU, input messages, memory) as steady only once its first iteration with input has completed, so that the balancer reads windowed metrics such as cpu_usage_10m after the job's initialization (state download, index build) has decayed out of them; the counters themselves are never reset. A job that has received no input for a minute is reported as steady as well, so an idle job reports its near-zero rates instead of staying unmeasured. Off by default while the change is rolled out pipeline by pipeline; temporary switch.
|
NYT::NFlow::TDynamicKeyVisitorStreamSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
period
|
Type: TDuration
Default value: 1d
|
|
buffer_row_limit
|
Type: NYT::NYTree::TSize
Default value: 5K
|
|
max_scan_rows_per_iteration
|
Type: NYT::NYTree::TSize
Default value: 10K
|
|
background_fill_period
|
Type: TDuration
Default value: 500ms
|
|
catchup_lag_threshold
|
Type: TDuration
Default value: 1m
|
|
catchup_speedup_multiplier
|
Type: double
Default value: 1.2
|
|
finite
|
Type: bool
Default value: true
Whether the visitor is finite: with %true it finishes once the streams it follows complete (see upstream_streams in the static spec), with %false it never finishes.
|
|
full_final_pass
|
Type: bool
Default value: true
|
NYT::NFlow::TDynamicMessageDistributorSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
send_queue_max_rows_per_batch
|
Type: NYT::NYTree::TSize
Default value: 10K
|
|
send_queue_max_bytes_per_batch
|
Type: NYT::NYTree::TSize
Default value: 10Mi
|
|
send_queue_batch_duration
|
Type: TDuration
Default value: 100ms
|
|
push_messages_timeout
|
Type: TDuration
Default value: 10s
|
|
compression_codec
|
Type: NYT::NCompression::ECodec
Default value: lz4
|
|
thread_count
|
Type: int
Default value: 4
|
|
hung_task_threshold
|
Type: TDuration
Default value: 5m
|
Additional parameters
|
max_processed_batch_size
|
Type: long
Default value: 500000
Maximum number of completed messages that a worker reports in a single PushMessages response.
|
NYT::NFlow::TDynamicOutputStoreSpec
Source: yt/yt/flow/library/cpp/common/spec.h
NYT::NFlow::TDynamicPartitionTracerSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
epochs_in_trace_context
|
Type: int
Default value: 3
|
|
trace_probability
|
Type: double
Default value: 0.0
|
|
trace_probability_partition_override
|
Type: THashMap<NYT::TStrongTypedef<NYT::TGuid, NYT::NFlow::TPartitionIdTag, NYT::TStrongTypedefOptions{true}>, double>
Default value: {}
|
|
wall_time_half_decay_period
|
Type: TDuration
Default value: 1m
|
NYT::NFlow::TDynamicPipelineSpec
Source: yt/yt/flow/library/cpp/common/spec.h
NYT::NFlow::TDynamicResourceSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
parameters
|
Type: NYT::TIntrusivePtr<NYT::NYTree::IMapNode>
Default value: {}
Dynamic parameters of the user resource class.
|
|
file_providers
|
Type: THashMap<NYT::TStrongTypedef<std::string, NYT::NFlow::TFileProviderIdTag, NYT::TStrongTypedefOptions{true}>, NYT::TIntrusivePtr<NYT::NFlow::TDynamicFileProviderSpec>>
Default value: {}
Dynamic parameters of the resource's named file providers. Changing the parameters immediately schedules discovery of the corresponding provider.
|
|
file_provider_discover_period
|
Type: TDuration
Default value: 30s
Discovery period for all named file providers of the resource.
|
|
file_provider_update_retry_period
|
Type: TDuration
Default value: 1m
Retry period after a file snapshot fails to download, initialize, or validate.
|
|
file_snapshot_min_creation_period
|
Type: TDuration
Default value: 5m
Minimum period between new complete file snapshots. More frequent provider updates are accumulated until the next snapshot is allowed.
|
|
file_snapshot_catalog_max_entries
|
Type: long
Default value: 1024
Maximum number of file snapshots retained in controller state. The current Active and Preparing snapshots are never evicted.
|
|
file_snapshot_rollout_warning_period
|
Type: TDuration
Default value: 15m
Time after publishing an Active file snapshot before an incomplete worker rollout is reported in the resource status.
|
NYT::NFlow::TDynamicRetryableRequestSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
min_inner_timeout
|
Type: TDuration
Default value: 1m
|
|
timeout
|
Type: TDuration
Default value: 5m
|
|
backoff
|
Type: NYT::TExponentialBackoffOptions
Default value:
{
"backoff_jitter" = 0.1;
"backoff_multiplier" = 2.;
"invocation_count" = 30;
"max_backoff" = 60000;
"min_backoff" = 1000;
}
|
|
lease_check_period
|
Type: TDuration
Default value: 1m
|
NYT::NFlow::TDynamicSimpleExternalStateJoinerSpec
Source: yt/yt/flow/library/cpp/computation/simple_external_state_manager.h
NYT::NFlow::TDynamicSimpleExternalStateManagerSpec
Source: yt/yt/flow/library/cpp/computation/simple_external_state_manager.h
This structure has no main parameters.
NYT::NFlow::TDynamicSinkSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
parameters
|
Type: NYT::TIntrusivePtr<NYT::NYTree::IMapNode>
Default value: {}
Arbitrary dynamic parameters of the corresponding Sink class.
|
NYT::NFlow::TDynamicSourceSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
parameters
|
Type: NYT::TIntrusivePtr<NYT::NYTree::IMapNode>
Default value: {}
Arbitrary dynamic parameters of the corresponding Source class.
|
|
draining
|
Type: bool
Default value: false
|
NYT::NFlow::TDynamicStateCacheSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
compressed_cache_weight
|
Type: NYT::NYTree::TSize
Default value: 1Gi
Limit for compressed.
|
|
uncompressed_cache_weight
|
Type: NYT::NYTree::TSize
Default value: 100Mi
Limit for uncompressed.
|
Source: yt/yt/flow/library/cpp/common/spec.h
NYT::NFlow::TDynamicStateJoinerSpec
Source: yt/yt/flow/library/cpp/common/spec.h
NYT::NFlow::TDynamicStateManagerSpec
Source: yt/yt/flow/library/cpp/common/spec.h
NYT::NFlow::TDynamicStaticTableKeyVisitorJoinerSpec
Source: yt/yt/flow/library/cpp/computation/static_table_key_visitor_joiner.h
|
Parameter
|
Description
|
|
read_attempts
|
Type: int
Default value: 3
Budget of attempts for a single source read. A read that exhausts the budget is considered unsuccessful and is handled according to unavailable_source_policy.
|
|
unavailable_source_backoff
|
Type: TDuration
Default value: 5m
How long an unsuccessful read marks the source as unavailable.
While the marker is active, there are no accesses to the source: each read immediately resolves according to unavailable_source_policy.
The first read after the window expires retries the source.
|
NYT::NFlow::TDynamicTableRequestSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
select_min_limit
|
Type: NYT::NYTree::TSize
Default value: 10
|
|
select_limit_multiplier
|
Type: long
Default value: 5
|
|
select_max_limit
|
Type: NYT::NYTree::TSize
Default value: 10K
|
NYT::NFlow::TDynamicThrottlerClassSpec
Source: yt/yt/flow/library/cpp/common/spec.h
Configuration of one distributed-throttler quota class.
|
Parameter
|
Description
|
|
weight
|
Type: double
Default value: 1.0
Positive finite weight defining the class's long-run share among active classes.
|
NYT::NFlow::TDynamicThrottlerSpec
Source: yt/yt/flow/library/cpp/common/spec.h
Configuration of a named distributed throttler.
|
Parameter
|
Description
|
|
limit
|
Type: std::optional<double>
Average quota issuance rate in units per second. none means unlimited.
|
|
period
|
Type: TDuration
Default value: 10s
Averaging period for the average rate in the token bucket algorithm: over period, the bucket is replenished to the volume of limit * period units. The same volume is the maximum possible burst.
|
|
request_period
|
Type: TDuration
Default value: 5s
Target interval between client requests to the controller (the prefetch size adjusts to this interval).
|
|
retrying_channel
|
Type: NYT::TIntrusivePtr<NYT::NRpc::TRetryingChannelConfig>
Default value:
{
"enable_exponential_retry_backoffs" = %true;
"retry_attempts" = 100;
"retry_timeout" = 600000;
}
Retry parameters for requests to the throttler service. The defaults are designed to survive a controller leader change even if it takes a long time.
|
|
rpc_timeout
|
Type: TDuration
Default value: 30s
Timeout for a single RequestQuota request.
|
|
classes
|
Type: THashMap<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TQuotaClassIdTag>, NYT::TIntrusivePtr<NYT::NFlow::TDynamicThrottlerClassSpec>>
Default value: {}
Named weighted quota classes. Keys must match [0-9A-Za-z_-]+. Backlogged classes share bandwidth in proportion to their weights and idle shares are redistributed.
|
|
max_grant_amount
|
Type: std::optional<long>
Maximum server scheduling chunk in absolute quota units. It bounds the delay before active classes are reconsidered. If unset, a single request is granted whole, so it holds the token bucket for its entire prefetch window and delays every other class by that long.
|
|
use_class_weights_as_limit
|
Type: bool
Default value: false
Reads the class weights as absolute rates: the issuance rate becomes the sum of the declared class weights, so a backlogged class is served at its weight in units per second. Requires at least one class and is mutually exclusive with limit. The reserved default class is not part of the sum.
|
NYT::NFlow::TDynamicTimerStoreSpec
Source: yt/yt/flow/library/cpp/common/spec.h
NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::NCompanion::TSwiftMapCompanionComputation>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
desired_partition_count
|
Type: std::optional<int>
|
|
min_partition_count
|
Type: std::optional<int>
|
|
max_partition_count
|
Type: std::optional<int>
|
|
sink_channel_multiplier
|
Type: std::optional<int>
|
|
desired_average_partition_cpu_load
|
Type: std::optional<double>
|
|
desired_average_partition_memory_used
|
Type: std::optional<double>
|
|
desired_average_partition_messages_per_second
|
Type: std::optional<double>
|
|
desired_average_partition_bytes_per_second
|
Type: std::optional<double>
|
|
desired_average_partition_timer_count
|
Type: std::optional<double>
|
|
allowed_partition_count_deviation
|
Type: std::optional<double>
|
|
partition_count_double_delay
|
Type: std::optional<TDuration>
|
|
partition_count_half_delay
|
Type: std::optional<TDuration>
|
|
weight_multiplier
|
Type: double
Default value: 1.0
|
|
interrupting_weight_multiplier
|
Type: double
Default value: 0.1
|
NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::NCompanion::TSwiftOrderedSourceCompanionComputation>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
desired_partition_count
|
Type: std::optional<int>
|
|
min_partition_count
|
Type: std::optional<int>
|
|
max_partition_count
|
Type: std::optional<int>
|
|
sink_channel_multiplier
|
Type: std::optional<int>
|
|
desired_average_partition_cpu_load
|
Type: std::optional<double>
|
|
desired_average_partition_memory_used
|
Type: std::optional<double>
|
|
desired_average_partition_messages_per_second
|
Type: std::optional<double>
|
|
desired_average_partition_bytes_per_second
|
Type: std::optional<double>
|
|
desired_average_partition_timer_count
|
Type: std::optional<double>
|
|
allowed_partition_count_deviation
|
Type: std::optional<double>
|
|
partition_count_double_delay
|
Type: std::optional<TDuration>
|
|
partition_count_half_delay
|
Type: std::optional<TDuration>
|
|
weight_multiplier
|
Type: double
Default value: 1.0
|
|
interrupting_weight_multiplier
|
Type: double
Default value: 0.1
|
|
max_read_window
|
Type: TDuration
Default value: 10m
|
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
desired_partition_count
|
Type: std::optional<int>
|
|
min_partition_count
|
Type: std::optional<int>
|
|
max_partition_count
|
Type: std::optional<int>
|
|
sink_channel_multiplier
|
Type: std::optional<int>
|
|
desired_average_partition_cpu_load
|
Type: std::optional<double>
|
|
desired_average_partition_memory_used
|
Type: std::optional<double>
|
|
desired_average_partition_messages_per_second
|
Type: std::optional<double>
|
|
desired_average_partition_bytes_per_second
|
Type: std::optional<double>
|
|
desired_average_partition_timer_count
|
Type: std::optional<double>
|
|
allowed_partition_count_deviation
|
Type: std::optional<double>
|
|
partition_count_double_delay
|
Type: std::optional<TDuration>
|
|
partition_count_half_delay
|
Type: std::optional<TDuration>
|
|
weight_multiplier
|
Type: double
Default value: 1.0
|
|
interrupting_weight_multiplier
|
Type: double
Default value: 0.1
|
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
desired_partition_count
|
Type: std::optional<int>
|
|
min_partition_count
|
Type: std::optional<int>
|
|
max_partition_count
|
Type: std::optional<int>
|
|
sink_channel_multiplier
|
Type: std::optional<int>
|
|
desired_average_partition_cpu_load
|
Type: std::optional<double>
|
|
desired_average_partition_memory_used
|
Type: std::optional<double>
|
|
desired_average_partition_messages_per_second
|
Type: std::optional<double>
|
|
desired_average_partition_bytes_per_second
|
Type: std::optional<double>
|
|
desired_average_partition_timer_count
|
Type: std::optional<double>
|
|
allowed_partition_count_deviation
|
Type: std::optional<double>
|
|
partition_count_double_delay
|
Type: std::optional<TDuration>
|
|
partition_count_half_delay
|
Type: std::optional<TDuration>
|
|
weight_multiplier
|
Type: double
Default value: 1.0
|
|
interrupting_weight_multiplier
|
Type: double
Default value: 0.1
|
NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::NSortedDynamicTable::TAsyncSink>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
max_rows_per_batch
|
Type: long
Default value: 1000
Maximum number of Flow messages in a batch: the batching sink flushes the accumulated batch when it reaches this threshold. The valid range is 1 through 2000.
|
|
max_bytes_per_batch
|
Type: long
Default value: 5242880
Maximum total size of Flow messages in a batch: the batching sink flushes the accumulated batch when it reaches this threshold. The valid range is 1 byte through 10 MiB.
|
Additional parameters
|
backoff_duration
|
Type: TDuration
Default value: 3s
Initial retry delay before jitter. Retry delays use exponential backoff with jitter; the delay before jitter is capped at one minute (or at this value when it is greater). The sink keeps retrying until the write succeeds or the job is cancelled, and reports the latest error in its status.
|
NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::NSortedDynamicTable::TSyncSink>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
This structure has no main parameters.
NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::NStaticTableConnector::TArrivalOrderTableSink>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
max_row_count
|
Type: long
Default value: 10000
Maximum number of messages in a table. Reaching the limit advances the sink to the next grid step, including future steps.
|
|
max_data_weight
|
Type: long
Default value: 1073741824
Maximum total message weight in a table. Reaching the limit advances the sink to the next grid step.
|
Additional parameters
|
transaction_timeout
|
Type: TDuration
Default value: 5m
Timeout of the master transaction that atomically creates a table and advances delivery progress.
|
|
retry_backoff
|
Type: TDuration
Default value: 1s
Delay between transaction retries. Retries continue until success or job cancellation.
|
NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::NStaticTableConnector::TSource>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
desired_table_process_time
|
Type: TDuration
Default value: 1h
The time to read one input table. The main parameter for controlling read speed.
|
|
max_rows_per_second
|
Type: double
Default value: 1000000000000.0
|
|
max_bytes_per_second
|
Type: double
Default value: 1000000000000.0
|
|
max_event_timestamp
|
Type: std::optional<unsigned long>
|
|
restart_instant
|
Type: TInstant
Default value: 1970-01-01T00:00:00.000000Z
Set this parameter to the current time in ISO 8601 to forget the current progress and start reading static tables from scratch. You can safely decrease the parameter value without triggering an additional read restart. A restart happens only if the current restart_instant is greater than the last restart_instant that the source saved internally. This parameter does not correlate with static table timestamps and does not filter them for the read restart.
|
|
allow_v1_migration
|
Type: bool
Default value: true
Controls departure from V1-compatible ordering. The default is true. A V1-shaped checkpoint is detected as V1, and a source with this flag enabled advances it one way through V1 → Draining → V2, or directly to V2 when no V1 table is in progress. Set the flag to false to postpone the initial cutover. Persisted Draining and V2 are authoritative and are never downgraded by changing this flag to false.
|
Additional parameters
These parameters are for fine tuning. It is not recommended to change them without a deep understanding of the system.
|
unavailable_threshold
|
Type: TDuration
Default value: 5m
How long the source must be continuously unavailable for the partition to count as stably unavailable. Only time during which the job was running and seeing the failure counts: a gap between restarts is not charged, and any successful answer from the source resets the total.
|
|
min_event_timestamp
|
Type: std::optional<unsigned long>
Tables whose EventTimestamp is less than MinEventTimestamp are not processed.
|
|
max_partition_count
|
Type: NYT::NYTree::TSize
Default value: 10K
Not recommended to change. An upper limit on the number of concurrently active partitions.
|
|
throttler_period
|
Type: TDuration
Default value: 10s
Not recommended to change. Adjusts the window in which the read speed from a partition is controlled.
|
|
desired_partition_process_time
|
Type: TDuration
Default value: 10m
Not recommended to change. Controls how large partitions the controller will create.
|
|
desired_partition_rows_per_second
|
Type: double
Default value: 1000.0
Not recommended to change. Controls how large partitions the controller will create.
|
|
desired_partition_bytes_per_second
|
Type: double
Default value: 1000000.0
Not recommended to change. Controls how large partitions the controller will create.
|
|
read_timeout
|
Type: TDuration
Default value: 5m
Not recommended to change. Recreates the table reader if it returns an empty response within the period.
|
NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::NStaticTableConnectorV2::TSource>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
desired_table_process_time
|
Type: TDuration
Default value: 1h
The time to read one input table. The main parameter for controlling read speed.
|
|
max_rows_per_second
|
Type: double
Default value: 1000000000000.0
|
|
max_bytes_per_second
|
Type: double
Default value: 1000000000000.0
|
|
max_event_timestamp
|
Type: std::optional<unsigned long>
|
|
restart_instant
|
Type: TInstant
Default value: 1970-01-01T00:00:00.000000Z
Set this parameter to the current time in ISO 8601 to forget the current progress and start reading static tables from scratch. You can safely decrease the parameter value without triggering an additional read restart. A restart happens only if the current restart_instant is greater than the last restart_instant that the source saved internally. This parameter does not correlate with static table timestamps and does not filter them for the read restart.
|
|
allow_v1_migration
|
Type: bool
Default value: true
Controls departure from V1-compatible ordering. The default is true. A V1-shaped checkpoint is detected as V1, and a source with this flag enabled advances it one way through V1 → Draining → V2, or directly to V2 when no V1 table is in progress. Set the flag to false to postpone the initial cutover. Persisted Draining and V2 are authoritative and are never downgraded by changing this flag to false.
|
Additional parameters
|
unavailable_threshold
|
Type: TDuration
Default value: 5m
How long the source must be continuously unavailable for the partition to count as stably unavailable. Only time during which the job was running and seeing the failure counts: a gap between restarts is not charged, and any successful answer from the source resets the total.
|
|
min_event_timestamp
|
Type: std::optional<unsigned long>
Tables whose EventTimestamp is less than MinEventTimestamp are not processed.
|
|
max_partition_count
|
Type: NYT::NYTree::TSize
Default value: 10K
Not recommended to change. An upper limit on the number of concurrently active partitions.
|
|
throttler_period
|
Type: TDuration
Default value: 10s
Not recommended to change. Adjusts the window in which the read speed from a partition is controlled.
|
|
desired_partition_process_time
|
Type: TDuration
Default value: 10m
Not recommended to change. Controls how large partitions the controller will create.
|
|
desired_partition_rows_per_second
|
Type: double
Default value: 1000.0
Not recommended to change. Controls how large partitions the controller will create.
|
|
desired_partition_bytes_per_second
|
Type: double
Default value: 1000000.0
Not recommended to change. Controls how large partitions the controller will create.
|
|
read_timeout
|
Type: TDuration
Default value: 5m
Not recommended to change. Recreates the table reader if it returns an empty response within the period.
|
NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::TAsyncHttpSink>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
at_most_once_strategy
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TAtMostOnceStrategyDynamicParameters>
Default value: {}
Dynamic parameters for at_most_once_strategy. Connector support varies; check the connector documentation before configuring them.
|
|
request_timeout
|
Type: TDuration
Default value: 1m
Total deadline for one message, including all HTTP attempts and retry delays. Exhausting it terminally stops the sink for the current job.
|
|
attempt_timeout
|
Type: TDuration
Default value: 10s
Maximum duration of one HTTP attempt. Each attempt is additionally limited by the time remaining before request_timeout; this value must not exceed request_timeout.
|
|
retry_initial_delay
|
Type: TDuration
Default value: 1s
Delay before the first retry. It must be between retry_minimum_delay and retry_maximum_delay.
|
|
retry_minimum_delay
|
Type: TDuration
Default value: 100ms
Lower bound for a retry delay after jitter is applied. It must not exceed retry_initial_delay.
|
|
retry_multiplier
|
Type: double
Default value: 2.0
Multiplier for exponential retry delay growth. Allowed values range from 1.01 through 100.0.
|
|
retry_maximum_delay
|
Type: TDuration
Default value: 30s
Upper bound for retry delays. It must not be less than retry_initial_delay.
|
|
retry_jitter_ratio
|
Type: double
Default value: 0.2
Random relative deviation applied to a retry delay. Allowed values range from 0.0 through 1.0.
|
|
max_attempt_count
|
Type: int
Default value: 5
Maximum number of HTTP attempts, including the initial attempt. Exhausting it terminally stops the sink for the current job. Allowed values range from 1 through 100.
|
NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::TAsyncMultiClusterQueueSink>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
at_most_once_strategy
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TAtMostOnceStrategyDynamicParameters>
Default value: {}
Dynamic parameters for at_most_once_strategy. Connector support varies; check the connector documentation before configuring them.
|
Additional parameters
|
flow_queue_meta_heartbeat_period
|
Type: TDuration
Default value: 10s
How often the controller will write heartbeats with meta containing watermarks to all queue partitions.
|
|
write_period
|
Type: TDuration
Default value: 100ms
|
|
max_rows_per_write
|
Type: long
Default value: 1000
|
|
max_bytes_per_write
|
Type: long
Default value: 1048576
|
|
backoff_duration
|
Type: TDuration
Default value: 1s
|
NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::TAsyncQueueSink>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
at_most_once_strategy
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TAtMostOnceStrategyDynamicParameters>
Default value: {}
Dynamic parameters for at_most_once_strategy. Connector support varies; check the connector documentation before configuring them.
|
Additional parameters
|
flow_queue_meta_heartbeat_period
|
Type: TDuration
Default value: 10s
How often the controller will write heartbeats with meta containing watermarks to all queue partitions.
|
|
write_period
|
Type: TDuration
Default value: 100ms
|
|
max_rows_per_write
|
Type: long
Default value: 1000
|
|
max_bytes_per_write
|
Type: long
Default value: 1048576
|
|
backoff_duration
|
Type: TDuration
Default value: 1s
|
NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::TAtLeastOnceClickHouseSink>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
write_timeout
|
Type: TDuration
Default value: 1m
Timeout of a single network operation (send/recv).
|
|
retry_backoff
|
Type: TDuration
Default value: 1s
Pause between insert retries.
|
|
max_insert_attempts
|
Type: long
Default value: 10
Maximum number of insert attempts for errors that cannot be reliably classified as transient (server errors, for example); once exhausted, the write fails. Known transient errors (network, protocol) are retried without a limit, and unrecoverable ones (validation) are not retried.
|
|
async_insert
|
Type: bool
Default value: false
Enables server-side asynchronous inserts. All sink classes emit async_insert=1 and wait_for_async_insert=1. The batching exactly-once sinks additionally provide their deduplication token and emit async_insert_deduplicate=1; the at-least-once and at-most-once sinks intentionally omit the deduplication token and the async_insert_deduplicate setting.
|
|
replay_horizon
|
Type: TDuration
Default value: 1d
Upper bound on replay lag, used to check the deduplication window. When the writer session starts before the first insert, it is compared with the server-side replicated_deduplication_window_seconds, or replicated_deduplication_window_seconds_for_async_inserts when async_insert is enabled. If the selected window is shorter, a warning is written to the worker log.
|
NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::TAtMostOnceClickHouseSink>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
write_timeout
|
Type: TDuration
Default value: 1m
Timeout of a single network operation (send/recv).
|
|
retry_backoff
|
Type: TDuration
Default value: 1s
Pause between insert retries.
|
|
max_insert_attempts
|
Type: long
Default value: 10
Maximum number of insert attempts for errors that cannot be reliably classified as transient (server errors, for example); once exhausted, the write fails. Known transient errors (network, protocol) are retried without a limit, and unrecoverable ones (validation) are not retried.
|
|
async_insert
|
Type: bool
Default value: false
Enables server-side asynchronous inserts. All sink classes emit async_insert=1 and wait_for_async_insert=1. The batching exactly-once sinks additionally provide their deduplication token and emit async_insert_deduplicate=1; the at-least-once and at-most-once sinks intentionally omit the deduplication token and the async_insert_deduplicate setting.
|
|
replay_horizon
|
Type: TDuration
Default value: 1d
Upper bound on replay lag, used to check the deduplication window. When the writer session starts before the first insert, it is compared with the server-side replicated_deduplication_window_seconds, or replicated_deduplication_window_seconds_for_async_inserts when async_insert is enabled. If the selected window is shorter, a warning is written to the worker log.
|
|
at_most_once_strategy
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TAtMostOnceStrategyDynamicParameters>
Default value: {}
Dynamic parameters for at_most_once_strategy. Connector support varies; check the connector documentation before configuring them.
|
NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::TClickHouseBatchingSink>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
write_timeout
|
Type: TDuration
Default value: 1m
Timeout of a single network operation (send/recv).
|
|
retry_backoff
|
Type: TDuration
Default value: 1s
Pause between insert retries.
|
|
max_insert_attempts
|
Type: long
Default value: 10
Maximum number of insert attempts for errors that cannot be reliably classified as transient (server errors, for example); once exhausted, the write fails. Known transient errors (network, protocol) are retried without a limit, and unrecoverable ones (validation) are not retried.
|
|
async_insert
|
Type: bool
Default value: false
Enables server-side asynchronous inserts. All sink classes emit async_insert=1 and wait_for_async_insert=1. The batching exactly-once sinks additionally provide their deduplication token and emit async_insert_deduplicate=1; the at-least-once and at-most-once sinks intentionally omit the deduplication token and the async_insert_deduplicate setting.
|
|
replay_horizon
|
Type: TDuration
Default value: 1d
Upper bound on replay lag, used to check the deduplication window. When the writer session starts before the first insert, it is compared with the server-side replicated_deduplication_window_seconds, or replicated_deduplication_window_seconds_for_async_inserts when async_insert is enabled. If the selected window is shorter, a warning is written to the worker log.
|
|
max_rows_per_batch
|
Type: long
Default value: 1000
Maximum number of Flow messages in a batch: the batching sink flushes the accumulated batch when it reaches this threshold. The valid range is 1 through 2000.
|
|
max_bytes_per_batch
|
Type: long
Default value: 5242880
Maximum total size of Flow messages in a batch: the batching sink flushes the accumulated batch when it reaches this threshold. The valid range is 1 byte through 10 MiB.
|
NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::TPassthroughComputation>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
desired_partition_count
|
Type: std::optional<int>
|
|
min_partition_count
|
Type: std::optional<int>
|
|
max_partition_count
|
Type: std::optional<int>
|
|
sink_channel_multiplier
|
Type: std::optional<int>
|
|
desired_average_partition_cpu_load
|
Type: std::optional<double>
|
|
desired_average_partition_memory_used
|
Type: std::optional<double>
|
|
desired_average_partition_messages_per_second
|
Type: std::optional<double>
|
|
desired_average_partition_bytes_per_second
|
Type: std::optional<double>
|
|
desired_average_partition_timer_count
|
Type: std::optional<double>
|
|
allowed_partition_count_deviation
|
Type: std::optional<double>
|
|
partition_count_double_delay
|
Type: std::optional<TDuration>
|
|
partition_count_half_delay
|
Type: std::optional<TDuration>
|
|
weight_multiplier
|
Type: double
Default value: 1.0
|
|
interrupting_weight_multiplier
|
Type: double
Default value: 0.1
|
NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::TQueueSource>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
pull_queue_timeout
|
Type: TDuration
Default value: 1m
Timeout for a read request from the queue.
|
Additional parameters
|
unavailable_threshold
|
Type: TDuration
Default value: 5m
How long the source must be continuously unavailable for the partition to count as stably unavailable. Only time during which the job was running and seeing the failure counts: a gap between restarts is not charged, and any successful answer from the source resets the total.
|
NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::TRandomSource>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
partition_count
|
Type: int
Default value: 3
|
|
partition_message_count
|
Type: std::optional<int>
|
|
message_size_mean
|
Type: int
Default value: 1024
|
|
message_count_mean
|
Type: int
Default value: 1000000
|
|
message_key_range
|
Type: int
Default value: 1024
|
|
reported_backlog_bytes_per_second
|
Type: std::optional<double>
|
|
reported_backlog_messages_per_second
|
Type: std::optional<double>
|
Additional parameters
|
unavailable_threshold
|
Type: TDuration
Default value: 5m
How long the source must be continuously unavailable for the partition to count as stably unavailable. Only time during which the job was running and seeing the failure counts: a gap between restarts is not charged, and any successful answer from the source resets the total.
|
NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::TServiceLogSource>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
desired_partition_count
|
Type: long
Default value: 5
Desired number of partitions. Estimate as 1 partition per 1 MB/s of stream.
|
|
desired_cycle_time
|
Type: TDuration
Default value: 12h
The source will try to maintain a message generation rate that traverses the entire table once every desired_cycle_time.
|
Additional parameters
|
unavailable_threshold
|
Type: TDuration
Default value: 5m
How long the source must be continuously unavailable for the partition to count as stably unavailable. Only time during which the job was running and seeing the failure counts: a gap between restarts is not charged, and any successful answer from the source resets the total.
|
|
throttler_period
|
Type: TDuration
Default value: 10s
Throttling period.
|
NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::TShardedClickHouseBatchingSink>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
write_timeout
|
Type: TDuration
Default value: 1m
Timeout of a single network operation (send/recv).
|
|
retry_backoff
|
Type: TDuration
Default value: 1s
Pause between insert retries.
|
|
max_insert_attempts
|
Type: long
Default value: 10
Maximum number of insert attempts for errors that cannot be reliably classified as transient (server errors, for example); once exhausted, the write fails. Known transient errors (network, protocol) are retried without a limit, and unrecoverable ones (validation) are not retried.
|
|
async_insert
|
Type: bool
Default value: false
Enables server-side asynchronous inserts. All sink classes emit async_insert=1 and wait_for_async_insert=1. The batching exactly-once sinks additionally provide their deduplication token and emit async_insert_deduplicate=1; the at-least-once and at-most-once sinks intentionally omit the deduplication token and the async_insert_deduplicate setting.
|
|
replay_horizon
|
Type: TDuration
Default value: 1d
Upper bound on replay lag, used to check the deduplication window. When the writer session starts before the first insert, it is compared with the server-side replicated_deduplication_window_seconds, or replicated_deduplication_window_seconds_for_async_inserts when async_insert is enabled. If the selected window is shorter, a warning is written to the worker log.
|
|
max_rows_per_batch
|
Type: long
Default value: 1000
Maximum number of Flow messages in a batch: the batching sink flushes the accumulated batch when it reaches this threshold. The valid range is 1 through 2000.
|
|
max_bytes_per_batch
|
Type: long
Default value: 5242880
Maximum total size of Flow messages in a batch: the batching sink flushes the accumulated batch when it reaches this threshold. The valid range is 1 byte through 10 MiB.
|
NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::TSwiftPassthroughComputation>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
desired_partition_count
|
Type: std::optional<int>
|
|
min_partition_count
|
Type: std::optional<int>
|
|
max_partition_count
|
Type: std::optional<int>
|
|
sink_channel_multiplier
|
Type: std::optional<int>
|
|
desired_average_partition_cpu_load
|
Type: std::optional<double>
|
|
desired_average_partition_memory_used
|
Type: std::optional<double>
|
|
desired_average_partition_messages_per_second
|
Type: std::optional<double>
|
|
desired_average_partition_bytes_per_second
|
Type: std::optional<double>
|
|
desired_average_partition_timer_count
|
Type: std::optional<double>
|
|
allowed_partition_count_deviation
|
Type: std::optional<double>
|
|
partition_count_double_delay
|
Type: std::optional<TDuration>
|
|
partition_count_half_delay
|
Type: std::optional<TDuration>
|
|
weight_multiplier
|
Type: double
Default value: 1.0
|
|
interrupting_weight_multiplier
|
Type: double
Default value: 0.1
|
NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::TSwiftPassthroughOrderedSourceComputation>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
desired_partition_count
|
Type: std::optional<int>
|
|
min_partition_count
|
Type: std::optional<int>
|
|
max_partition_count
|
Type: std::optional<int>
|
|
sink_channel_multiplier
|
Type: std::optional<int>
|
|
desired_average_partition_cpu_load
|
Type: std::optional<double>
|
|
desired_average_partition_memory_used
|
Type: std::optional<double>
|
|
desired_average_partition_messages_per_second
|
Type: std::optional<double>
|
|
desired_average_partition_bytes_per_second
|
Type: std::optional<double>
|
|
desired_average_partition_timer_count
|
Type: std::optional<double>
|
|
allowed_partition_count_deviation
|
Type: std::optional<double>
|
|
partition_count_double_delay
|
Type: std::optional<TDuration>
|
|
partition_count_half_delay
|
Type: std::optional<TDuration>
|
|
weight_multiplier
|
Type: double
Default value: 1.0
|
|
interrupting_weight_multiplier
|
Type: double
Default value: 0.1
|
|
max_read_window
|
Type: TDuration
Default value: 10m
|
NYT::NFlow::TDynamicUnitedParameters<NYT::NFlow::TSyncQueueSink>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
This structure has no main parameters.
Additional parameters
|
flow_queue_meta_heartbeat_period
|
Type: TDuration
Default value: 10s
How often the controller will write heartbeats with meta containing watermarks to all queue partitions.
|
NYT::NFlow::TEvenLoadThresholds
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
spread
|
Type: std::optional<double>
Minimal gap between the maximal and the minimal per-worker usage of the resource (in the resource's own units: cores for CPU, bytes for memory) from which the load counts as uneven.
|
|
ratio
|
Type: std::optional<double>
Minimal ratio of the maximal per-worker usage to the minimal one from which the load counts as uneven. At least 1.
|
NYT::NFlow::TEventTimestampAssignerSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
column
|
Type: std::optional<std::string>
The time column.
|
|
format
|
Type: NYT::NFlow::ETimestampFormat
Default value: seconds
The time format used (seconds and milliseconds).
|
|
limit_by_system_timestamp
|
Type: bool
Default value: false
Limit EventTimestamp to the SystemTimestamp value to protect against corrupted values from the future.
|
NYT::NFlow::TExternalStateJoinerSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
external_state_joiner_class_name
|
Type: std::string
Default value: NYT::NFlow::TSimpleExternalStateJoiner
Full class name of the external state joiner. Must be registered with the YT_FLOW_DEFINE_EXTERNAL_STATE_JOINER macro (or be a library implementation such as NYT::NFlow::TSimpleExternalStateJoiner).
|
|
client_provider_resource_id
|
Type: std::optional<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TResourceIdTag>>
If set, the framework searches for a static resource with this name (must implement IYTClientProvider) and uses its Get() as the joiner's YT client. In this case, the cluster from the table's rich path is ignored. Mutually exclusive with client_factory_resource_id (this parameter takes priority).
|
|
client_factory_resource_id
|
Type: std::optional<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TResourceIdTag>>
If set, the framework searches for a static resource with this name (must implement IYTClientProvider) and uses GetClient(cluster) to get a YT client for the cluster from the table's rich path. Mutually exclusive with client_provider_resource_id.
|
|
join_on
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TStateJoinSpec>
Default value: {}
Description of the mapping between the computation input stream keys and the joiner keys (see TStateJoinSpec).
|
|
auto_preload
|
Type: bool
Default value: true
If true (default), the framework automatically calls PreloadKeyStates on this joiner before each DoProcess, forming keys using join_on. If false, the computation is responsible for calling Client.PreloadKeyStates(IInputContextPtr) or PreloadKeyStates(THashSet<TKey>) before GetState.
|
|
parameters
|
Type: NYT::TIntrusivePtr<NYT::NYTree::IMapNode>
Parameters of the specific joiner implementation.
|
NYT::NFlow::TExternalStateManagerSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
external_state_manager_class_name
|
Type: std::string
Default value: NYT::NFlow::TSimpleExternalStateManager
Full class name of the external state manager. Must be registered with the YT_FLOW_DEFINE_EXTERNAL_STATE_MANAGER macro (or be a library implementation such as NYT::NFlow::TSimpleExternalStateManager).
|
|
auto_preload
|
Type: bool
Default value: true
If true (default), the framework automatically calls PreloadKeyStates on this manager before each DoProcess with every message, timer and visit key of the epoch. If false, the computation is responsible for calling Client.PreloadKeyStates(IInputContextPtr), PreloadKeyStates(IInputContextPtr, TExtractKeysOptions) or PreloadKeyStates(THashSet<TKey>) before GetState; GetState on a key that was not preloaded throws. Not allowed for companion computations.
|
|
parameters
|
Type: NYT::TIntrusivePtr<NYT::NYTree::IMapNode>
Parameters of the selected external state manager implementation. The expected schema is determined by class_name.
|
NYT::NFlow::TFetcherInJoinerSpec
Source: yt/yt/flow/library/cpp/connectors/servicelog/joiner.h
|
Parameter
|
Description
|
|
table_path
|
Type: NYT::NYPath::TRichYPath
Default value: TString("")
Path to the table including the cluster (multiple clusters can be specified).
|
|
value_columns
|
Type: std::optional<THashSet<std::string>>
Non-key columns to read (by default, all columns are read).
|
|
attempts
|
Type: long
Default value: 1
|
|
retry_timeout
|
Type: TDuration
Default value: 5s
|
|
fetch_type
|
Type: NYT::NFlow::EFetchType
Default value: table_reader
Important parameter that affects the type of load on YT. See the EFetchType description.
|
|
prefix
|
Type: std::string
Default value: TString("")
Column name prefix to be added for non-key table columns.
|
NYT::NFlow::TFileProviderSpec
Source: yt/yt/flow/library/cpp/common/file_provider.h
|
Parameter
|
Description
|
|
file_provider_class_name
|
Type: std::string
Required parameter
Name of an IFileProvider implementation registered with YT_FLOW_DEFINE_FILE_PROVIDER.
|
|
parameters
|
Type: NYT::TIntrusivePtr<NYT::NYTree::IMapNode>
Default value: {}
Parameters of the registered file-provider implementation.
|
|
postprocess_command
|
Type: std::optional<std::string>
Optional shell command that transforms the downloaded tree before resource initialization. The successful result is cached by provider revision and exact command.
|
|
postprocess_timeout
|
Type: TDuration
Default value: 1m
Maximum command duration. Changing the timeout does not invalidate an already cached postprocessing result.
|
NYT::NFlow::TFlowNodeConfig
Source: yt/yt/flow/library/cpp/runner/config.h
|
Parameter
|
Description
|
|
logging
|
Type: NYT::TIntrusivePtr<NYT::NLogging::TLogManagerConfig>
Default value:
{
"high_backlog_watermark" = 100000;
"low_backlog_watermark" = 100000;
"min_disk_space" = 0;
"rules" = [
{
"exclude_categories" = [];
"max_level" = "maximum";
"min_level" = "info";
"writers" = [
"Stderr";
];
};
];
"writers" = {
"Stderr" = {
"common_fields" = {};
"enable_host_field" = %false;
"enable_native_tags" = %false;
"enable_source_location" = %false;
"enable_system_fields" = %true;
"format" = "plain_text";
"type" = "stderr";
"yson_format" = "text";
};
};
}
Logging settings.
|
|
jaeger
|
Type: NYT::TIntrusivePtr<NYT::NTracing::TJaegerTracerConfig>
Default value: {}
|
|
pipe_io_dispatcher
|
Type: NYT::TIntrusivePtr<NYT::NPipeIO::TPipeIODispatcherConfig>
Default value: {}
|
|
solomon_registry
|
Type: NYT::TIntrusivePtr<NYT::NProfiling::TSolomonRegistryConfig>
Default value: {}
|
|
error_backtrace_enricher
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TBacktraceEnricherSpec>
Default value: {}
|
|
query_engine_config
|
Type: NYT::TIntrusivePtr<NYT::NQueryClient::TQueryEngineConfig>
Default value: {}
|
|
cluster_url
|
Type: std::string
Required parameter
YTsaurus cluster to use.
|
|
path
|
Type: TString
Required parameter
Working directory of the pipeline on the cluster_url cluster.
|
|
proxy_role
|
Type: std::optional<std::string>
rpc proxy role.
|
|
rpc_port
|
Type: int
Default value: 0
Main port for communication between workers and controllers, can be different for each instance.
|
|
monitoring_port
|
Type: int
Default value: 0
API for obtaining metrics.
|
|
companion
|
Type: NYT::TIntrusivePtr<NYT::NFlow::NCompanion::TCompanionConfig>
Companion process parameters. Needed by any worker that runs a companion (Python, Java, Go, C++ companion). Filled in automatically for a vanilla launch.
|
|
controller
|
Type: NYT::TIntrusivePtr<NYT::NFlow::NController::TControllerConfig>
Default value: {}
Controller parameters.
|
|
worker
|
Type: NYT::TIntrusivePtr<NYT::NFlow::NWorker::TWorkerConfig>
Default value: {}
Worker parameters.
|
|
authenticator
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TAuthenticatorConfig>
Default value: {}
|
|
bus_server
|
Type: NYT::TIntrusivePtr<NYT::NBus::NTcp::TBusServerConfig>
Default value: {}
|
|
rpc_server
|
Type: NYT::TIntrusivePtr<NYT::NRpc::TServerConfig>
Default value: {}
|
|
core_dumper
|
Type: NYT::TIntrusivePtr<NYT::NCoreDump::TCoreDumperConfig>
|
|
solomon_exporter
|
Type: NYT::TIntrusivePtr<NYT::NProfiling::TSolomonExporterConfig>
Default value: {'enable_solomon_aggregates': true}
|
|
solomon_proxy
|
Type: NYT::TIntrusivePtr<NYT::NProfiling::TSolomonProxyConfig>
Default value: {}
|
|
http_client_config
|
Type: NYT::TIntrusivePtr<NYT::NHttp::TClientConfig>
Default value: {}
|
|
https_client_config
|
Type: NYT::TIntrusivePtr<NYT::NHttps::TClientConfig>
Default value: {}
|
|
http_poller_threads
|
Type: int
Default value: 1
|
|
abort_on_unrecognized_options
|
Type: bool
Default value: true
When set to true, Flow Node will not start if there are unknown options in the config.
|
|
ignore_singletons_dynamic_config
|
Type: bool
Default value: false
|
|
enable_porto_resource_tracker
|
Type: bool
Default value: true
|
Additional parameters
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
window
|
Type: TDuration
Default value: 5m
Time window.
|
|
threshold
|
Type: double
Default value: 0.01
Keys that appear in the stream more often than threshold are considered heavy hitters. The algorithm uses O(1 / threshold) memory.
|
|
limit
|
Type: long
Default value: 5
Maximum number of heavy hitter keys to show.
|
NYT::NFlow::TIdlePartitionsSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
duration
|
Type: TDuration
Default value: 1m
The time a partition must be empty before it is considered to have no writes.
|
|
max_ratio
|
Type: double
Default value: 0.4
The proportion of partitions that can be excluded from the watermark calculation by this logic. The conceptual purpose of this option is to protect against watermark advancement when the data provider stops writing messages to the queue due to an incident on their side. Setting the option to 1.0 effectively disables this protection.
|
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
time_type
|
Type: NYT::NFlow::ETimeType
Default value: event_time
Time for sorting: event_time, system_time, real_time
|
|
stream_delays
|
Type: THashMap<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TStreamIdTag>, TDuration>
Default value: {}
Delays for input streams. If a stream is not specified, there is no delay
|
NYT::NFlow::TKeyStateRow
Source: yt/yt/flow/library/cpp/controller/state_access.h
A single row in the key_states, external_key_states, joined_external_key_states sections of the read-states response.
|
Parameter
|
Description
|
|
computation_id
|
Type: NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TComputationIdTag>
Required parameter
ID of the Computation to which the key belongs.
|
|
key
|
Type: NYT::TStrongTypedef<NYT::NFlow::TCompactUnversionedOwningRow, NYT::NFlow::TKeyTag, NYT::TStrongTypedefOptions{true}>
Required parameter
Key value (YSON Dict {column = value; ...} by key_schema).
|
|
states
|
Type: THashMap<std::string, NYT::NYson::TYsonString>
Required parameter
Content of the states under this key: a dict state_name → YSON value. The YSON value here is the serialized user payload of the state (the type is determined by the specific State/ExternalState in the Computation code; for built-in types like counters it is just a number, for arbitrary user structs — a YSON dict with all their fields).
For the external_key_states / joined_external_key_states sections, the state name matches the name of the external-state client (what the process function passes to IRuntimeInitContext::InitExternalStateClient, for example /state).
|
NYT::NFlow::TKeyVisitorStreamSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
names
|
Type: std::optional<THashSet<std::string>>
|
|
external_names
|
Type: std::optional<THashSet<std::string>>
|
|
upstream_streams
|
Type: std::optional<THashSet<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TStreamIdTag>>>
List of the computation's input and source streams whose completion the visitor waits for. Only meaningful with finite=%true (see the dynamic spec): with finite=%false the visitor never finishes, so the parameter has no effect.
When unset, the visitor waits for every input and source stream, that is, it finishes together with the whole input of the computation.
Narrowing the list matters for a cyclic topology: when the process function emits from ProcessVisit into an output that comes back to the same computation as an input stream. Such a stream only ends once the visitor stops scanning the state, so it cannot be followed.
In a production pipeline this makes no difference: sources are infinite, input streams never complete, and no value of upstream_streams starts a final pass. Narrowing is needed where all pipeline sources are finite and the pipeline must reach completed, that is, in integration tests.
|
|
bucket_count
|
Type: int
Default value: 8
|
NYT::NFlow::TLateDataPartitionsSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
value
|
Type: NYT::TStrongTypedef<double, NYT::NFlow::TWatermarkPercentileTag, NYT::TStrongTypedefOptions{true}>
Default value: 100.0
|
|
delay
|
Type: TDuration
Default value: 1m
|
NYT::NFlow::TLoadThroughputThrottlerSpec
Source: yt/yt/flow/library/cpp/misc/load_throughput_throttler.h
|
Parameter
|
Description
|
|
limit
|
Type: std::optional<double>
|
|
period
|
Type: TDuration
Default value: 1s
|
|
alpha
|
Type: double
Default value: 0.999
|
|
initial_row_size
|
Type: double
Default value: 1000.0
|
|
initial_key_size
|
Type: double
Default value: 500.0
|
NYT::NFlow::TMatchedStates
Source: yt/yt/flow/library/cpp/controller/state_access.h
Breakdown of matched rows by category. Used in TDeleteStatesResponse.
NYT::NFlow::TMatchedStatesBucket
Source: yt/yt/flow/library/cpp/controller/state_access.h
Count of matched rows for one of the categories of the delete-states response.
|
Parameter
|
Description
|
|
total
|
Type: long
Default value: 0
Total number of matched rows in this category.
|
|
details
|
Type: THashMap<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TComputationIdTag>, THashMap<std::string, long>>
Default value: {}
Breakdown by computation_id → state_name → number of rows.
|
NYT::NFlow::TMessageSerializer
Source: yt/yt/flow/library/cpp/common/message-inl.h
|
Parameter
|
Description
|
|
message_id
|
Type: NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TMessageIdTag>
Required parameter
Unique message ID.
|
|
system_timestamp
|
Type: NYT::TStrongTypedef<unsigned long, NYT::NFlow::TSystemTimestampTag, NYT::TStrongTypedefOptions{true}>
Required parameter
Timestamp of the specific message creation.
|
|
alignment_timestamp
|
Type: NYT::TStrongTypedef<unsigned long, NYT::NFlow::TSystemTimestampTag, NYT::TStrongTypedefOptions{true}>
Required parameter
|
|
event_timestamp
|
Type: NYT::TStrongTypedef<unsigned long, NYT::NFlow::TSystemTimestampTag, NYT::TStrongTypedefOptions{true}>
Required parameter
Timestamp of the real event associated with this message.
|
|
stream_id
|
Type: NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TStreamIdTag>
Required parameter
The stream this message belongs to.
|
|
payload
|
Type: NYT::TStrongTypedef<NYT::NFlow::TCompactUnversionedOwningRow, NYT::NFlow::TPayloadTag, NYT::TStrongTypedefOptions{true}>
Required parameter
The data itself.
|
|
payload_schema
|
Type: NYT::TIntrusivePtr<NYT::NTableClient::TTableSchema>
Required parameter
Payload schema, populated based on stream_id.
|
NYT::NFlow::TPartitionStateRow
Source: yt/yt/flow/library/cpp/controller/state_access.h
A single row in the partition_states section of the read-states response.
|
Parameter
|
Description
|
|
computation_id
|
Type: std::optional<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TComputationIdTag>>
ID of the Computation. May be absent if the partition is not bound to a specific Computation.
|
|
partition_id
|
Type: NYT::TStrongTypedef<NYT::TGuid, NYT::NFlow::TPartitionIdTag, NYT::TStrongTypedefOptions{true}>
Required parameter
Partition ID.
|
|
states
|
Type: THashMap<std::string, NYT::NYson::TYsonString>
Required parameter
Content of the partition states: a map state_name → YSON value.
|
NYT::NFlow::TPipelineSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
validate_binary_version
|
Type: bool
Default value: false
|
|
binary_version
|
Type: std::string
Default value: TString("")
|
|
validate_binary_checksum
|
Type: bool
Default value: false
|
|
binary_checksum
|
Type: std::string
Default value: TString("")
|
|
computations
|
Type: THashMap<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TComputationIdTag>, NYT::TIntrusivePtr<NYT::NFlow::TComputationSpec>>
Default value: {}
A named enumeration of all pipeline nodes. Keys must match [0-9A-Za-z_-]+.
|
|
resources
|
Type: THashMap<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TResourceIdTag>, NYT::TIntrusivePtr<NYT::NFlow::TResourceSpec>>
Default value: {}
A named enumeration of all resources. Keys must match [0-9A-Za-z_-]+.
|
|
streams
|
Type: THashMap<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TStreamIdTag>, NYT::TIntrusivePtr<NYT::NFlow::TStreamSpec>>
Default value: {}
A named enumeration of all data streams. Keys must match [0-9A-Za-z_-]+.
|
NYT::NFlow::TPivotFinderSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
part_count
|
Type: long
Default value: 16
|
|
window_size
|
Type: long
Default value: 16384
|
NYT::NFlow::TReadStatesArg
Source: yt/yt/flow/library/cpp/controller/state_access.h
Argument of the read-states command.
The command supports three usage modes (determined by the set of filled fields):
- Only
computation_id — all states for all partitions of the specified Computation.
computation_id + key — point lookup by key in key_states (and in external states) of the specified Computation.
- Only
partition_id — all partition_states of the specified partition. Additionally, if this partition belongs to SourceComputation (each such partition is responsible for exactly one SourceKey, unlike regular partitions with a [LowerKey; UpperKey) range), key_states under this SourceKey are also loaded immediately.
|
Parameter
|
Description
|
|
computation_id
|
Type: std::optional<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TComputationIdTag>>
ID of the Computation. Required for modes 1 and 2 (see the command description).
|
|
partition_id
|
Type: std::optional<NYT::TStrongTypedef<NYT::TGuid, NYT::NFlow::TPartitionIdTag, NYT::TStrongTypedefOptions{true}>>
Partition ID. Required for mode 3 (see the command description).
|
|
key
|
Type: std::optional<NYT::NYson::TYsonString>
Key for point lookup (mode 2). Accepted in two forms:
- YSON Dict
{column = value; ...} — for each target table (key_states of the Computation and each table of the external_state_manager/joiner with its own key_schema_override) the dict is decomposed by its key schema; columns absent from the schema are ignored.
- YSON List of positional values — passed to tables as-is, without adaptation to their schema.
|
|
name
|
Type: std::optional<std::string>
Filter by exact state name; applies to all sections (key/partition/external). If set, only rows with this state name are returned in each section.
|
|
target
|
Type: NYT::NFlow::EFlowStateTarget
Default value: all
Which categories of states to request. Default: all.
|
|
limit
|
Type: long
Default value: 10
Maximum number of rows returned in each response section independently.
Sections refer to the response fields: key_states, partition_states, external_key_states, joined_external_key_states. Each has its own budget of limit rows.
Within the external_key_states / joined_external_key_states sections, the budget is fairly divided among the active external states of this Computation.
|
NYT::NFlow::TReadStatesResponse
Source: yt/yt/flow/library/cpp/controller/state_access.h
Response of the read-states command. Each section is populated independently according to its own limit.
|
Parameter
|
Description
|
|
key_states
|
Type: std::vector<NYT::NFlow::TKeyStateRow>
Default value: []
States bound to keys (by the Computation's key schema).
|
|
partition_states
|
Type: std::vector<NYT::NFlow::TPartitionStateRow>
Default value: []
States bound to partitions.
|
|
external_key_states
|
Type: std::vector<NYT::NFlow::TKeyStateRow>
Default value: []
States from external_state_managers (mutable, can be deleted via delete-states).
|
|
joined_external_key_states
|
Type: std::vector<NYT::NFlow::TKeyStateRow>
Default value: []
States observed via external_state_joiner (read-only, join sources).
|
|
errors
|
Type: std::vector<std::string>
Default value: []
Errors that occurred when forming the response to the request itself (for example, failed to adapt the key to the schema of a specific external state or a List/Lookup to its table failed). These are not related to errors in the pipeline's state operation. If such an error occurs for one external state, the other sections are still populated successfully.
|
NYT::NFlow::TResourceDescription
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
alias
|
Type: std::optional<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TResourceIdTag>>
|
|
worker
|
Type: bool
Default value: true
|
|
controller
|
Type: bool
Default value: true
|
NYT::NFlow::TResourceSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
resource_class_name
|
Type: std::string
Required parameter
Name of the corresponding class.
Must be registered using the YT_FLOW_DEFINE_RESOURCE macro.
|
|
parameters
|
Type: NYT::TIntrusivePtr<NYT::NYTree::IMapNode>
Default value: {}
Arbitrary parameters for the corresponding class.
|
|
file_providers
|
Type: THashMap<NYT::TStrongTypedef<std::string, NYT::NFlow::TFileProviderIdTag, NYT::TStrongTypedefOptions{true}>, NYT::TIntrusivePtr<NYT::NFlow::TFileProviderSpec>>
Default value: {}
Named file providers of the resource. The controller attaches a file snapshot to a target revision only after discovering every name; a controller-specific target spec may be published independently. Workers materialize the exact revisions from that target. Such resources are worker-only and must have controller = %false on every required_resource_ids path that can reach them. Keys must match [0-9A-Za-z_-]+.
|
|
dependencies
|
Type: THashMap<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TResourceIdTag>, NYT::TIntrusivePtr<NYT::NFlow::TResourceDescription>>
Default value: {}
Resources that this resource depends on. Mirrors the structure of required_resource_ids from Computation.
|
|
required_capabilities
|
Type: THashMap<std::string, long>
Default value: {}
|
|
preload_required
|
Type: bool
Default value: false
|
|
always_on
|
Type: bool
Default value: false
If true, the resource is loaded in advance when the unit starts and stays in memory for its entire lifetime (never unloaded) instead of being loaded lazily on first access. It is loaded only on units where at least one computation requires it (see worker / controller in required_resource_ids): a resource needed only on a worker is not loaded on a controller, and a resource not required by any computation is not loaded at all. Mutually exclusive with preload_required.
|
NYT::NFlow::TSimpleExternalStateJoinerSpec
Source: yt/yt/flow/library/cpp/computation/simple_external_state_manager.h
|
Parameter
|
Description
|
|
path
|
Type: NYT::NYPath::TRichYPath
Required parameter
Path to the dynamic table with states (including the cluster in a rich path).
|
NYT::NFlow::TSimpleExternalStateManagerSpec
Source: yt/yt/flow/library/cpp/computation/simple_external_state_manager.h
|
Parameter
|
Description
|
|
path
|
Type: NYT::NYPath::TRichYPath
Required parameter
Path to the dynamic table with states (including the cluster in a rich path).
|
NYT::NFlow::TSimpleRunnerConfig
Source: yt/yt/flow/library/cpp/runner/simple_runner_program.h
|
Parameter
|
Description
|
|
logging
|
Type: NYT::TIntrusivePtr<NYT::NLogging::TLogManagerConfig>
Default value:
{
"high_backlog_watermark" = 100000;
"low_backlog_watermark" = 100000;
"min_disk_space" = 0;
"rules" = [
{
"exclude_categories" = [];
"max_level" = "maximum";
"min_level" = "info";
"writers" = [
"Stderr";
];
};
];
"writers" = {
"Stderr" = {
"common_fields" = {};
"enable_host_field" = %false;
"enable_native_tags" = %false;
"enable_source_location" = %false;
"enable_system_fields" = %true;
"format" = "plain_text";
"type" = "stderr";
"yson_format" = "text";
};
};
}
Logging settings.
|
|
jaeger
|
Type: NYT::TIntrusivePtr<NYT::NTracing::TJaegerTracerConfig>
Default value: {}
|
|
pipe_io_dispatcher
|
Type: NYT::TIntrusivePtr<NYT::NPipeIO::TPipeIODispatcherConfig>
Default value: {}
|
|
solomon_registry
|
Type: NYT::TIntrusivePtr<NYT::NProfiling::TSolomonRegistryConfig>
Default value: {}
|
|
error_backtrace_enricher
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TBacktraceEnricherSpec>
Default value: {}
|
|
query_engine_config
|
Type: NYT::TIntrusivePtr<NYT::NQueryClient::TQueryEngineConfig>
Default value: {}
|
|
cluster_url
|
Type: std::string
Required parameter
Name of the cluster where the pipeline is located. The tail after the last / is treated as a proxy role and takes precedence over proxy_role — e.g. hahn/flow means cluster hahn with role flow.
|
|
proxy_role
|
Type: std::optional<std::string>
Required parameter
The proxy role to use for performing operations on the pipeline.
|
|
path
|
Type: TString
Required parameter
Path to the pipeline on YTsaurus.
|
|
spec
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TPipelineSpec>
Required parameter
Static spec with which to start the pipeline.
|
|
dynamic_spec
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TDynamicPipelineSpec>
Default value: {}
Dynamic spec with which to start the pipeline.
|
|
abort_on_unrecognized_options
|
Type: bool
Default value: false
|
|
abort_on_specs_parseability_error
|
Type: bool
Default value: false
|
|
vanilla
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TVanillaConfig>
If set and enable=%true, the runner starts a vanilla operation on YTsaurus and runs the federation inside it, instead of starting the controller/workers locally.
|
|
set_flow_core_target
|
Type: bool
Default value: true
|
|
direct_controller_commands
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TDirectControllerCommandsConfig>
Default value: {}
|
Additional parameters
NYT::NFlow::TSinkSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
sink_class_name
|
Type: std::string
Default value: TString("")
Name of the corresponding class.
Must be registered using YT_FLOW_DEFINE_SINK.
|
|
input_stream_ids
|
Type: THashSet<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TStreamIdTag>>
Default value: []
List of input streams.
Only streams specified in output_stream_ids of the entire Computation are allowed here.
|
|
parameters
|
Type: NYT::TIntrusivePtr<NYT::NYTree::IMapNode>
Default value: {}
Arbitrary parameters for the corresponding class.
|
NYT::NFlow::TSourceSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
source_class_name
|
Type: std::string
Default value: TString("")
Name of the corresponding class.
Must be registered using YT_FLOW_DEFINE_SOURCE.
|
|
parameters
|
Type: NYT::TIntrusivePtr<NYT::NYTree::IMapNode>
Default value: {}
Arbitrary parameters for the corresponding class.
|
NYT::NFlow::TStateJoinSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
key_schema_override
|
Type: NYT::TIntrusivePtr<NYT::NTableClient::TTableSchema>
If set, the schema is used as the joiner key instead of the computation's group_by_schema.
|
|
key_provider_streams
|
Type: std::optional<THashSet<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TStreamIdTag>>>
If set, the joiner keys are taken only from messages and timers with the specified stream_id; nullopt means "from all computation input streams".
|
NYT::NFlow::TStateJoinerSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
computation_id
|
Type: NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TComputationIdTag>
Required parameter
The Computation whose internal state is read by this joiner.
|
|
state_name
|
Type: std::string
Required parameter
The state client name of the target Computation — the prefix passed to InitClient (starts with /).
|
|
join_on
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TStateJoinSpec>
Default value: {}
Description of the mapping between the computation input stream keys and the target Computation's state key (see TStateJoinSpec). key_schema_override must match the target Computation's group_by_schema.
|
|
auto_preload
|
Type: bool
Default value: true
If true (default), the framework automatically calls PreloadKeyStates on this joiner before each DoProcess, forming keys using join_on. If false, the computation is responsible for calling PreloadKeyStates before GetState.
|
NYT::NFlow::TStaticTableKeyVisitorJoinerSpec
Source: yt/yt/flow/library/cpp/computation/static_table_key_visitor_joiner.h
|
Parameter
|
Description
|
|
path
|
Type: NYT::NYPath::TRichYPath
Required parameter
Path to a static sorted source table (including the cluster in a rich path). Its key column prefix must match the computation's group_by_schema by names and types. The computed (partition) column must be materialized with the values of its expression.
|
|
unavailable_source_policy
|
Type: NYT::NFlow::EUnavailableSourcePolicy
Default value: retry
Behavior when reading the source fails (exhausting the read_attempts budget).
retry (default) — the error drops the traversal iteration and reading is retried. The traversal does not advance past the unread range.
mark_unreadable — the error is swallowed, the traversal continues, and keys in the unread range are resolved with an uninitialized accessor (IsInitialized() == false).
|
NYT::NFlow::TStreamSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
class_name
|
Type: std::optional<std::string>
|
|
schema
|
Type: NYT::TIntrusivePtr<NYT::NTableClient::TTableSchema>
Required parameter
Stream schema. It is a dynamic table schema. Cannot contain sorting or computed columns.
We recommend making schema changes by completely stopping the pipeline. In this case, no output messages remain inside the pipeline.
|
|
migration_function
|
Type: std::string
Default value: NYT::NFlow::DefaultMigrationFunction
|
NYT::NFlow::TTableFetcherSpec
Source: yt/yt/flow/library/cpp/connectors/servicelog/fetcher.h
|
Parameter
|
Description
|
|
table_path
|
Type: NYT::NYPath::TRichYPath
Default value: TString("")
Path to the table including the cluster (multiple clusters can be specified).
|
|
value_columns
|
Type: std::optional<THashSet<std::string>>
Non-key columns to read (by default, all columns are read).
|
|
attempts
|
Type: long
Default value: 1
|
|
retry_timeout
|
Type: TDuration
Default value: 5s
|
|
fetch_type
|
Type: NYT::NFlow::EFetchType
Default value: table_reader
Important parameter that affects the type of load on YT. See the EFetchType description.
|
NYT::NFlow::TTableJoinerSpec
Source: yt/yt/flow/library/cpp/connectors/servicelog/joiner.h
NYT::NFlow::TTimerSerializer
Source: yt/yt/flow/library/cpp/common/timer-inl.h
|
Parameter
|
Description
|
|
message_id
|
Type: NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TMessageIdTag>
Required parameter
Unique timer ID.
|
|
system_timestamp
|
Type: NYT::TStrongTypedef<unsigned long, NYT::NFlow::TSystemTimestampTag, NYT::TStrongTypedefOptions{true}>
Required parameter
Timestamp of the timer creation.
|
|
event_timestamp
|
Type: NYT::TStrongTypedef<unsigned long, NYT::NFlow::TSystemTimestampTag, NYT::TStrongTypedefOptions{true}>
Required parameter
Timestamp of the real event associated with this timer.
|
|
stream_id
|
Type: NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TStreamIdTag>
Required parameter
The stream this message belongs to.
|
|
key
|
Type: NYT::TStrongTypedef<NYT::NFlow::TCompactUnversionedOwningRow, NYT::NFlow::TKeyTag, NYT::TStrongTypedefOptions{true}>
Required parameter
Key of the timer.
|
|
key_schema
|
Type: NYT::TIntrusivePtr<NYT::NTableClient::TTableSchema>
Required parameter
Key schema, matches the group_by_schema of the corresponding Computation.
|
|
trigger_timestamp
|
Type: NYT::TStrongTypedef<unsigned long, NYT::NFlow::TSystemTimestampTag, NYT::TStrongTypedefOptions{true}>
Required parameter
Time for the timer to fire.
|
NYT::NFlow::TTimerSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
time_type
|
Type: NYT::NFlow::ETimeType
Default value: event_time
The time type used by the timer: event_time, system_time, real_time
|
|
streams
|
Type: std::optional<THashSet<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TStreamIdTag>>>
List of streams that the timer monitors.
|
|
streams_with_delays
|
Type: std::optional<THashMap<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TStreamIdTag>, TDuration>>
List of streams that the timer monitors with an individual delay specified
|
|
deduplicate_equal_timestamps
|
Type: bool
Default value: true
Enables timer deduplication: if you try to create a timer with the same key and trigger_timestamp, only one timer remains — the one with the smaller EventTimestamp value
|
NYT::NFlow::TUnavailablePartitionGroupsSpec
Source: yt/yt/flow/library/cpp/common/spec.h
This structure has no main parameters.
NYT::NFlow::TUnitedParameters<NYT::NFlow::NCompanion::TSwiftMapCompanionComputation>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
allow_batching_with_relaxed_guarantees
|
Type: bool
Default value: false
Allows combining several input messages into one output message (batching). Useful for folding many small messages into one larger message to reduce the message count load on downstream partitions. Default is %false: each output message is produced by exactly one parent. MessageId is deterministically inherited from the parent, providing exactly-once semantics with deterministic user functions. When set to %true, an output message can have multiple parents. A parent is considered processed only when all its children are processed. On restarts, batch boundaries may differ, so the same parent message could end up in multiple different batched children — downstream computations must be prepared to see each parent's contents more than once (at-least-once semantics). Additionally, batched output messages receive a MessageId derived deterministically from the set of parent MessageIds: a replay of the same batch yields the same MessageId and is deduplicated, while a batch with a different composition yields a new one. This means that the order of messages by MessageId within a single key on the downstream Computation side can be disrupted. If subsequent pipeline stages rely on MessageId ordering within a key, you must rewrite that logic to account for this.
|
NYT::NFlow::TUnitedParameters<NYT::NFlow::NCompanion::TSwiftOrderedSourceCompanionComputation>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
This structure has no main parameters.
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
processing_mode
|
Type: NYT::NFlow::EProcessingMode
Default value: exactly_once
|
|
internal_states
|
Type: std::optional<THashSet<std::string>>
|
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
internal_states
|
Type: std::optional<THashSet<std::string>>
|
NYT::NFlow::TUnitedParameters<NYT::NFlow::NSortedDynamicTable::TAsyncSink>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
table_path
|
Type: NYT::NYPath::TRichYPath
Required parameter
Path to the target sorted dynamic table with the cluster specified.
|
|
column_filter
|
Type: std::optional<THashSet<std::string>>
Which columns from the message to write to the table (default — all). When using the delete_rows parameter, you must specify the key columns of the target table, otherwise deletions will fail.
|
|
aggregate_columns
|
Type: std::optional<THashSet<std::string>>
Aggregate writes are unsupported: omit this parameter or provide an empty list. An asynchronous write may replay an already committed batch after recovery, which makes non-idempotent aggregation unsafe.
|
|
delete_rows
|
Type: bool
Default value: false
If set to true, rows will be deleted instead of inserted. In this mode, make sure that column_filter selects the key columns of the target table.
|
|
require_sync_replica
|
Type: bool
Default value: true
If set to false, the presence of a synchronous replica is not checked when writing rows to the target table. The default value is true.
|
Additional parameters
|
update_partition_count_period
|
Type: TDuration
Default value: 1m
How often the controller should update the number of tablets in the table (the controller changes the number of receiver channels accordingly). For a chaos replicated table, the controller uses the enabled data replica with the lexicographically smallest replica id; its current sync/async mode does not affect the choice. Resharding that replica may change the number of receiver channels.
|
NYT::NFlow::TUnitedParameters<NYT::NFlow::NSortedDynamicTable::TSyncSink>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
table_path
|
Type: NYT::NYPath::TRichYPath
Required parameter
Path to the target sorted dynamic table with the cluster specified.
|
|
column_filter
|
Type: std::optional<THashSet<std::string>>
Which columns from the message to write to the table (default — all). When using the delete_rows parameter, you must specify the key columns of the target table, otherwise deletions will fail.
|
|
aggregate_columns
|
Type: std::optional<THashSet<std::string>>
List of columns for which the aggregation function is considered during writing. By default, the value in the column is overwritten.
|
|
delete_rows
|
Type: bool
Default value: false
If set to true, rows will be deleted instead of inserted. In this mode, it is important that the key columns of the target table are selected in column_filter.
Note
A Sink can work in either insert or delete mode, but not simultaneously. If you need both insert and delete functionality, use StateManager.
|
|
require_sync_replica
|
Type: bool
Default value: true
If set to false, the presence of a synchronous replica is not checked when writing rows to the target table. The default value is true — the check is performed.
|
Additional parameters
|
update_partition_count_period
|
Type: TDuration
Default value: 1m
How often the controller should update the number of tablets in the table (the controller changes the number of receiver channels accordingly). For a chaos replicated table, the controller uses the enabled data replica with the lexicographically smallest replica id; its current sync/async mode does not affect the choice. Resharding that replica may change the number of receiver channels.
|
NYT::NFlow::TUnitedParameters<NYT::NFlow::NStaticTableConnector::TArrivalOrderTableSink>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
output_directory
|
Type: NYT::NYPath::TRichYPath
Required parameter
Directory for the output static tables. Its @progress attribute holds the delivery progress: the (pipeline, computation, sink id) owner, the shared table sequence and a per-partition frontier. Every sink needs its own directory: if it finds a foreign owner, the sink fails and asks for the attribute to be removed by hand.
|
|
table_period
|
Type: TDuration
Default value: 5m
Step of the gap-free table grid. An empty table for step T is created without a gap only when the input stream system watermark is known and strictly greater than T + table_period. Empty tables are never created ahead of the wall clock.
|
|
table_ttl
|
Type: TDuration
Required parameter
Output table lifetime: a table is removed table_ttl after its table_timestamp (via Cypress expiration_time). The table_ttl / table_period ratio is capped at 40000 by the Cypress directory child count limit.
|
|
table_name_format
|
Type: std::string
Default value: %Y-%m-%dT%H:%M:%SZ
Format of the UTC timestamp in a table name. It must render the timestamp losslessly (validated by a round-trip), or different slots would collide on one name.
|
|
data_weight_column
|
Type: std::optional<std::string>
Optional column with a user-defined message weight for the max_data_weight limit (int64 or uint64, non-negative). Without it, and for null values, the message byte size is used.
|
NYT::NFlow::TUnitedParameters<NYT::NFlow::NStaticTableConnector::TSource>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
finite
|
Type: bool
Default value: false
Treat the source as finite: remember the number of messages in it on startup and transition the stream to the completed state after all those messages are read.
|
|
tables
|
Type: std::optional<std::vector<NYT::NYPath::TRichYPath>>
List of tables to read with cluster specifications. This parameter is an alternative to tables_path. Each element must be a table: a symlink is not dereferenced to the target table, and any non-table node causes an error.
|
|
tables_path
|
Type: std::optional<NYT::NYPath::TRichYPath>
Path to a directory with static tables, including either one cluster or an ordered list of replica clusters. This parameter is an alternative to tables.
|
|
table_name_filter
|
Type: NYT::TIntrusivePtr<NYT::NRe2::TRe2>
Regular expression to filter tables by name: only tables whose name matches the expression are read. The expression syntax is RE2.
|
|
event_timestamp_locator
|
Type: NYT::TIntrusivePtr<NYT::NFlow::NStaticTableConnector::TTableTimestampLocatorSpec>
Default value: {'attribute': 'key'}
By default, takes the timestamp from the table name. This time corresponds to the data creation time and is forwarded to the EventTimestamp of messages. Tables with the same timestamp are ordered persistently; new tables must not appear behind the already processed event-time frontier.
|
|
system_timestamp_locator
|
Type: NYT::TIntrusivePtr<NYT::NFlow::NStaticTableConnector::TTableTimestampLocatorSpec>
Default value: {'attribute': 'creation_time'}
By default, takes the timestamp from the table creation time. This time corresponds to when the data is written to the source. That is, the moment when the pipeline can see this data and start reading. This time is forwarded to the SystemTimestamp of messages.
|
|
use_planned_timestamps
|
Type: bool
Default value: false
Use the persisted planned start time of each row range for both input message timestamps. The plan is fixed when the controller starts a table and survives restarts and rate changes. Timestamp locators still determine table discovery and ordering. Disabled by default.
|
|
ignore_symlinks
|
Type: bool
Default value: false
A flag that allows ignoring symlinks inside the table folder.
|
|
skip_non_table_nodes
|
Type: bool
Default value: false
Skip non-table nodes in the input directory instead of failing.
|
|
idle_watermark_delay
|
Type: std::optional<TDuration>
Default value: 3600000
The delay for advancing the watermark from the current time while the source is not reading a table. A table immediately advances the watermark to its event timestamp without subtracting this delay. Set this field to # to disable clock-based advancement; new tables still advance the watermark.
|
|
failover_delay
|
Type: TDuration
Default value: 5m
For replicated input, how long the active cluster must stay unavailable before the current table fails over.
|
Additional parameters
|
update_info_period
|
Type: TDuration
Default value: 15s
Period for refreshing the source's auxiliary partition information. A connector may use this tick for status requests and read-session liveness checks.
|
|
byte_size_alpha
|
Type: double
Default value: 0.05
Exponential smoothing coefficient for the average byte total and message count per offset: larger values make the estimate react to new data faster.
|
NYT::NFlow::TUnitedParameters<NYT::NFlow::NStaticTableConnectorV2::TSource>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
finite
|
Type: bool
Default value: false
Treat the source as finite: remember the number of messages in it on startup and transition the stream to the completed state after all those messages are read.
|
|
tables
|
Type: std::optional<std::vector<NYT::NYPath::TRichYPath>>
List of tables to read with cluster specifications. This parameter is an alternative to tables_path. Each element must be a table: a symlink is not dereferenced to the target table, and any non-table node causes an error.
|
|
tables_path
|
Type: std::optional<NYT::NYPath::TRichYPath>
Path to a directory with static tables, including either one cluster or an ordered list of replica clusters. This parameter is an alternative to tables.
|
|
table_name_filter
|
Type: NYT::TIntrusivePtr<NYT::NRe2::TRe2>
Regular expression to filter tables by name: only tables whose name matches the expression are read. The expression syntax is RE2.
|
|
event_timestamp_locator
|
Type: NYT::TIntrusivePtr<NYT::NFlow::NStaticTableConnector::TTableTimestampLocatorSpec>
Default value: {'attribute': 'key'}
By default, takes the timestamp from the table name. This time corresponds to the data creation time and is forwarded to the EventTimestamp of messages. Tables with the same timestamp are ordered persistently; new tables must not appear behind the already processed event-time frontier.
|
|
system_timestamp_locator
|
Type: NYT::TIntrusivePtr<NYT::NFlow::NStaticTableConnector::TTableTimestampLocatorSpec>
Default value: {'attribute': 'creation_time'}
By default, takes the timestamp from the table creation time. This time corresponds to when the data is written to the source. That is, the moment when the pipeline can see this data and start reading. This time is forwarded to the SystemTimestamp of messages.
|
|
use_planned_timestamps
|
Type: bool
Default value: false
Use the persisted planned start time of each row range for both input message timestamps. The plan is fixed when the controller starts a table and survives restarts and rate changes. Timestamp locators still determine table discovery and ordering. Disabled by default.
|
|
ignore_symlinks
|
Type: bool
Default value: false
A flag that allows ignoring symlinks inside the table folder.
|
|
skip_non_table_nodes
|
Type: bool
Default value: false
Skip non-table nodes in the input directory instead of failing.
|
|
idle_watermark_delay
|
Type: std::optional<TDuration>
Default value: 3600000
The delay for advancing the watermark from the current time while the source is not reading a table. A table immediately advances the watermark to its event timestamp without subtracting this delay. Set this field to # to disable clock-based advancement; new tables still advance the watermark.
|
|
failover_delay
|
Type: TDuration
Default value: 5m
For replicated input, how long the active cluster must stay unavailable before the current table fails over.
|
Additional parameters
|
update_info_period
|
Type: TDuration
Default value: 15s
Period for refreshing the source's auxiliary partition information. A connector may use this tick for status requests and read-session liveness checks.
|
|
byte_size_alpha
|
Type: double
Default value: 0.05
Exponential smoothing coefficient for the average byte total and message count per offset: larger values make the estimate react to new data faster.
|
NYT::NFlow::TUnitedParameters<NYT::NFlow::TAsyncHttpSink>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
at_most_once_strategy
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TAtMostOnceStrategyParameters>
Optional at-most-once delivery strategy. Connector support varies; check the connector documentation before enabling it.
|
|
url
|
Type: std::string
Required parameter
Target HTTP or HTTPS URL. The URL must contain a host.
|
|
payload_column
|
Type: std::string
Required parameter
Name of the payload column. The input stream must contain exactly one String column with this name. A non-null value becomes the request body byte-for-byte unchanged, whether it contains protobuf, YSON, JSON, or arbitrary String bytes. A null value is acknowledged and skipped without an HTTP request.
|
|
headers
|
Type: THashMap<std::string, std::string>
Default value: {}
Static HTTP request headers. Header names must be valid HTTP tokens. Set Content-Type here if the receiver uses it to interpret the opaque request body. Header values are stored in the pipeline Spec in Cypress and are readable by principals with Spec read access, including through spec dumps. There is no secret-safe header mechanism yet. Do not put OAuth tokens, API tokens, or other secrets here. Do not use this sink with endpoints that require secret headers unless access controls and an external mechanism make that safe.
|
|
idempotency_header
|
Type: std::string
Default value: Idempotency-Key
Name of the per-message idempotency header. Its value is the hexadecimal representation of the stable Flow message ID and remains unchanged across retries and replay. The default is Idempotency-Key. Set another valid name to rename the header or an empty string to disable it. A static header with the same name is rejected.
|
|
keep_alive
|
Type: bool
Default value: true
Whether to reuse HTTP connections. Every configured HTTP sink in every active partition job owns a distinct sink instance, client, and idle-connection pool. When disabled, the pool keeps no idle connections regardless of max_idle_connections.
|
|
max_redirect_count
|
Type: int
Default value: 0
Maximum number of HTTP redirects followed across all retry attempts for one delivery while preserving the POST method, body, and headers. The default is zero, so redirect following is opt-in. Trust every possible redirect destination.
|
|
max_idle_connections
|
Type: int
Default value: 8
Maximum number of idle connections retained by each HTTP or HTTPS client of every configured sink instance and active partition job when keep_alive is enabled. The default is eight per client, so a sink that follows redirects across both schemes can retain up to twice that number. An ordered-source partition corresponds to one source key, so there is no additional per-key multiplier. If several sinks target the same endpoint, size its capacity as the sum across all configured sinks and active partition jobs.
|
NYT::NFlow::TUnitedParameters<NYT::NFlow::TAsyncMultiClusterQueueSink>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
queue_path
|
Type: NYT::NYPath::TRichYPath
Required parameter
Path to the queue including the cluster.
|
|
write_flow_queue_meta
|
Type: bool
Default value: false
Whether to write meta to the queue (watermarks are written in meta).
|
|
flow_queue_meta_column
|
Type: std::string
Default value: flow_queue_meta
Column to which to write meta.
|
|
producer_path
|
Type: NYT::NYPath::TRichYPath
Required parameter
Path to the queue producer with the cluster specified.
|
|
require_sync_replica
|
Type: bool
Default value: true
The eponymous parameter when writing to the queue. Whether write to the queue without synchronous replicas is allowed.
|
|
tablet_index_expression
|
Type: std::optional<std::string>
Verbatim routing mode. A QL expression (builtins only, e.g. farm_hash) over the message columns whose value is written to the system column $tablet_index verbatim. Mutually exclusive with tablet_index_routing_hash_expression; when neither expression is set, routing is off ($tablet_index is not written and the tablet is chosen by the driver). Sharded queues are inconvenient to operate (object→tablet affinity is tightly coupled and a reshard remaps keys), so enable routing only when genuinely needed. Routing is supported only on sync queue sinks; async sinks reject the routing params at spec load (YTFLOW-766).
|
|
tablet_index_routing_hash_expression
|
Type: std::optional<std::string>
Hash routing mode. A QL expression (builtins only, e.g. farm_hash) yielding a uint64 hash that is reduced to a tablet index by tablet_index_routing_hash_policy over tablet_count. Mutually exclusive with tablet_index_expression. Supported only on sync queue sinks; async sinks reject the routing params (YTFLOW-766).
|
|
tablet_index_routing_hash_policy
|
Type: std::optional<NYT::NFlow::EQueueTabletIndexRoutingHashPolicy>
The policy for reducing the tablet_index_routing_hash_expression hash to a tablet index. Required together with it. range — contiguous equal-width hash ranges (rangeSize = 2^64 / tablet_count), recommended: a consumer range-partitioned by the same key reads only its own tablet. modulo — hash % tablet_count; discouraged: Flow computations are range-partitioned by key, so a modulo-partitioned queue forces every reader to read every tablet (a full mesh on read).
|
|
tablet_count
|
Type: std::optional<long>
The number of tablets used to reduce tablet_index_routing_hash_expression to a tablet index. Optional. When unset, the sink periodically re-reads the target queue's @tablet_count (from the master cache) and follows a reshard without a restart. When set explicitly, the value is fixed (no queue read) and a reshard is reflected only by changing the config.
|
|
at_most_once_strategy
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TAtMostOnceStrategyParameters>
Optional at-most-once delivery strategy. Connector support varies; check the connector documentation before enabling it.
|
|
column_filter
|
Type: std::optional<THashSet<std::string>>
Which columns from the message to write to the queue (default — all).
|
|
use_clusters
|
Type: bool
Default value: false
Whether to use <clusters=[...]> instead of <cluster=...> from the rich path in producer and queue
|
Additional parameters
|
update_partition_count_period
|
Type: TDuration
Default value: 1m
How often the controller should update the partition count in the queue (the controller changes the computation partition count according to the queue partition count).
|
|
update_partition_count_retry_min_backoff
|
Type: TDuration
Default value: 1s
Initial delay before retrying a failed partition count update, subject to jitter. The value is capped by update_partition_count_period; subsequent delays increase exponentially.
|
NYT::NFlow::TUnitedParameters<NYT::NFlow::TAsyncQueueSink>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
queue_path
|
Type: NYT::NYPath::TRichYPath
Required parameter
Path to the queue including the cluster.
|
|
write_flow_queue_meta
|
Type: bool
Default value: false
Whether to write meta to the queue (watermarks are written in meta).
|
|
flow_queue_meta_column
|
Type: std::string
Default value: flow_queue_meta
Column to which to write meta.
|
|
producer_path
|
Type: NYT::NYPath::TRichYPath
Required parameter
Path to the queue producer with the cluster specified.
|
|
require_sync_replica
|
Type: bool
Default value: true
The eponymous parameter when writing to the queue. Whether write to the queue without synchronous replicas is allowed.
|
|
tablet_index_expression
|
Type: std::optional<std::string>
Verbatim routing mode. A QL expression (builtins only, e.g. farm_hash) over the message columns whose value is written to the system column $tablet_index verbatim. Mutually exclusive with tablet_index_routing_hash_expression; when neither expression is set, routing is off ($tablet_index is not written and the tablet is chosen by the driver). Sharded queues are inconvenient to operate (object→tablet affinity is tightly coupled and a reshard remaps keys), so enable routing only when genuinely needed. Routing is supported only on sync queue sinks; async sinks reject the routing params at spec load (YTFLOW-766).
|
|
tablet_index_routing_hash_expression
|
Type: std::optional<std::string>
Hash routing mode. A QL expression (builtins only, e.g. farm_hash) yielding a uint64 hash that is reduced to a tablet index by tablet_index_routing_hash_policy over tablet_count. Mutually exclusive with tablet_index_expression. Supported only on sync queue sinks; async sinks reject the routing params (YTFLOW-766).
|
|
tablet_index_routing_hash_policy
|
Type: std::optional<NYT::NFlow::EQueueTabletIndexRoutingHashPolicy>
The policy for reducing the tablet_index_routing_hash_expression hash to a tablet index. Required together with it. range — contiguous equal-width hash ranges (rangeSize = 2^64 / tablet_count), recommended: a consumer range-partitioned by the same key reads only its own tablet. modulo — hash % tablet_count; discouraged: Flow computations are range-partitioned by key, so a modulo-partitioned queue forces every reader to read every tablet (a full mesh on read).
|
|
tablet_count
|
Type: std::optional<long>
The number of tablets used to reduce tablet_index_routing_hash_expression to a tablet index. Optional. When unset, the sink periodically re-reads the target queue's @tablet_count (from the master cache) and follows a reshard without a restart. When set explicitly, the value is fixed (no queue read) and a reshard is reflected only by changing the config.
|
|
at_most_once_strategy
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TAtMostOnceStrategyParameters>
Optional at-most-once delivery strategy. Connector support varies; check the connector documentation before enabling it.
|
|
column_filter
|
Type: std::optional<THashSet<std::string>>
Which columns from the message to write to the queue (default — all).
|
Additional parameters
|
update_partition_count_period
|
Type: TDuration
Default value: 1m
How often the controller should update the partition count in the queue (the controller changes the computation partition count according to the queue partition count).
|
|
update_partition_count_retry_min_backoff
|
Type: TDuration
Default value: 1s
Initial delay before retrying a failed partition count update, subject to jitter. The value is capped by update_partition_count_period; subsequent delays increase exponentially.
|
NYT::NFlow::TUnitedParameters<NYT::NFlow::TAtLeastOnceClickHouseSink>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
host
|
Type: std::string
Default value: TString("")
The single ClickHouse host of the unsharded form. For exactly-once, this form is used with TClickHouseBatchingSink. Required unless hosts or shard_hosts is set: exactly one of the three host forms must be present in the spec. The host must be reachable and expose a table with a supported engine and the complete ordered schema (name, type, default_kind, default_expression).
|
|
port
|
Type: unsigned short
Default value: 9000
Native TCP port. Shared by all shards.
|
|
hosts
|
Type: std::vector<std::string>
Default value: []
Flat list of hosts of one unsharded table, all sharing port. For exactly-once, this form is used with TClickHouseBatchingSink. Required unless host or shard_hosts is set. Every host must be reachable and match the engine and complete ordered schema (name, type, default_kind, default_expression). For Replicated*MergeTree, zookeeper_name and zookeeper_path must match, while replica_name may differ. Multiple hosts with plain MergeTree or SharedMergeTree are rejected; single-host SharedMergeTree remains supported. This target-identity rule applies to all delivery guarantees. A single-entry list is rejected — use host instead.
|
|
host_selection_policy
|
Type: NYT::NFlow::EClickHouseHostSelectionPolicy
Default value: ordered_round_robin
Endpoint order within each shard. The default ordered_round_robin preserves the configured order. random_start chooses a uniformly distributed starting position once when a client is constructed and rotates the list. clickhouse-cpp then exhausts all endpoints in round-robin order. The policy does not change row routing, deduplication tokens, or the topology fingerprint.
|
|
shard_hosts
|
Type: THashMap<std::string, std::vector<std::string>>
Default value: {}
Mapping from a shard name to the list of that shard's hosts, used by the multi-shard form. Required by TShardedClickHouseBatchingSink and rejected by TClickHouseBatchingSink; the other sink classes accept any of the three host forms. Every endpoint must be reachable. Within each shard, the engine and complete ordered schema (name, type, default_kind, default_expression) must match; schemas must also match across shards. For replicas, zookeeper_name and zookeeper_path must match, while replica_name may differ. Multiple hosts with plain MergeTree or SharedMergeTree are rejected; single-host SharedMergeTree remains supported. This target-identity rule applies to all delivery guarantees. Shard names must be non-empty and consist of letters, digits, _ and -; a shard's host list must not be empty. A shard name is a permanent identifier: it is part of the deduplication token.
|
|
sharding_key_columns
|
Type: std::vector<std::string>
Default value: []
Columns whose values produce the routing key. Requires shard_hosts: the unsharded forms reject this parameter during spec validation. An empty list routes by message id, that is, without co-location.
|
|
user
|
Type: std::string
Default value: default
ClickHouse user.
|
|
password_env_var
|
Type: std::string
Default value: TString("")
Name of the environment variable holding the password; resolved via GetEnv when the client is created. An empty value means connecting without a password.
|
|
database
|
Type: std::string
Default value: default
Database of the target table. Shared by all shards.
|
|
table
|
Type: std::string
Required parameter
Target table. Shared by all shards, so <database>.<table> must exist on every host of every shard.
|
|
codec
|
Type: NYT::NFlow::EClickHouseCodec
Default value: lz4
Native protocol compression.
|
|
enable_tls
|
Type: bool
Default value: false
Connect over TCP+TLS instead of plain TCP. The port is not changed implicitly: set port to the server's secure native port (usually 9440).
|
|
tls_ca_files
|
Type: std::vector<std::string>
Default value: []
Paths to root certificate (CA) files used to verify the server certificate. An empty list is valid when TLS is disabled; when TLS is enabled, it means that the system root certificates are used. A nonempty list requires enable_tls = true; otherwise configuration validation fails.
|
|
tls_ca_directory
|
Type: std::string
Default value: TString("")
Path to a directory of root certificates (CA) used to verify the server certificate. An empty value is valid when TLS is disabled. A nonempty path requires enable_tls = true; otherwise configuration validation fails.
|
|
tls_skip_verification
|
Type: bool
Default value: false
Skip TLS session verification (the server certificate and so on). Insecure; use only for test environments with self-signed certificates. Setting this option to true requires enable_tls = true; otherwise configuration validation fails. The default value, false, is valid when TLS is disabled.
|
NYT::NFlow::TUnitedParameters<NYT::NFlow::TAtMostOnceClickHouseSink>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
host
|
Type: std::string
Default value: TString("")
The single ClickHouse host of the unsharded form. For exactly-once, this form is used with TClickHouseBatchingSink. Required unless hosts or shard_hosts is set: exactly one of the three host forms must be present in the spec. The host must be reachable and expose a table with a supported engine and the complete ordered schema (name, type, default_kind, default_expression).
|
|
port
|
Type: unsigned short
Default value: 9000
Native TCP port. Shared by all shards.
|
|
hosts
|
Type: std::vector<std::string>
Default value: []
Flat list of hosts of one unsharded table, all sharing port. For exactly-once, this form is used with TClickHouseBatchingSink. Required unless host or shard_hosts is set. Every host must be reachable and match the engine and complete ordered schema (name, type, default_kind, default_expression). For Replicated*MergeTree, zookeeper_name and zookeeper_path must match, while replica_name may differ. Multiple hosts with plain MergeTree or SharedMergeTree are rejected; single-host SharedMergeTree remains supported. This target-identity rule applies to all delivery guarantees. A single-entry list is rejected — use host instead.
|
|
host_selection_policy
|
Type: NYT::NFlow::EClickHouseHostSelectionPolicy
Default value: ordered_round_robin
Endpoint order within each shard. The default ordered_round_robin preserves the configured order. random_start chooses a uniformly distributed starting position once when a client is constructed and rotates the list. clickhouse-cpp then exhausts all endpoints in round-robin order. The policy does not change row routing, deduplication tokens, or the topology fingerprint.
|
|
shard_hosts
|
Type: THashMap<std::string, std::vector<std::string>>
Default value: {}
Mapping from a shard name to the list of that shard's hosts, used by the multi-shard form. Required by TShardedClickHouseBatchingSink and rejected by TClickHouseBatchingSink; the other sink classes accept any of the three host forms. Every endpoint must be reachable. Within each shard, the engine and complete ordered schema (name, type, default_kind, default_expression) must match; schemas must also match across shards. For replicas, zookeeper_name and zookeeper_path must match, while replica_name may differ. Multiple hosts with plain MergeTree or SharedMergeTree are rejected; single-host SharedMergeTree remains supported. This target-identity rule applies to all delivery guarantees. Shard names must be non-empty and consist of letters, digits, _ and -; a shard's host list must not be empty. A shard name is a permanent identifier: it is part of the deduplication token.
|
|
sharding_key_columns
|
Type: std::vector<std::string>
Default value: []
Columns whose values produce the routing key. Requires shard_hosts: the unsharded forms reject this parameter during spec validation. An empty list routes by message id, that is, without co-location.
|
|
user
|
Type: std::string
Default value: default
ClickHouse user.
|
|
password_env_var
|
Type: std::string
Default value: TString("")
Name of the environment variable holding the password; resolved via GetEnv when the client is created. An empty value means connecting without a password.
|
|
database
|
Type: std::string
Default value: default
Database of the target table. Shared by all shards.
|
|
table
|
Type: std::string
Required parameter
Target table. Shared by all shards, so <database>.<table> must exist on every host of every shard.
|
|
codec
|
Type: NYT::NFlow::EClickHouseCodec
Default value: lz4
Native protocol compression.
|
|
enable_tls
|
Type: bool
Default value: false
Connect over TCP+TLS instead of plain TCP. The port is not changed implicitly: set port to the server's secure native port (usually 9440).
|
|
tls_ca_files
|
Type: std::vector<std::string>
Default value: []
Paths to root certificate (CA) files used to verify the server certificate. An empty list is valid when TLS is disabled; when TLS is enabled, it means that the system root certificates are used. A nonempty list requires enable_tls = true; otherwise configuration validation fails.
|
|
tls_ca_directory
|
Type: std::string
Default value: TString("")
Path to a directory of root certificates (CA) used to verify the server certificate. An empty value is valid when TLS is disabled. A nonempty path requires enable_tls = true; otherwise configuration validation fails.
|
|
tls_skip_verification
|
Type: bool
Default value: false
Skip TLS session verification (the server certificate and so on). Insecure; use only for test environments with self-signed certificates. Setting this option to true requires enable_tls = true; otherwise configuration validation fails. The default value, false, is valid when TLS is disabled.
|
|
at_most_once_strategy
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TAtMostOnceStrategyParameters>
Optional at-most-once delivery strategy. Connector support varies; check the connector documentation before enabling it.
|
NYT::NFlow::TUnitedParameters<NYT::NFlow::TClickHouseBatchingSink>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
host
|
Type: std::string
Default value: TString("")
The single ClickHouse host of the unsharded form. For exactly-once, this form is used with TClickHouseBatchingSink. Required unless hosts or shard_hosts is set: exactly one of the three host forms must be present in the spec. The host must be reachable and expose a table with a supported engine and the complete ordered schema (name, type, default_kind, default_expression).
|
|
port
|
Type: unsigned short
Default value: 9000
Native TCP port. Shared by all shards.
|
|
hosts
|
Type: std::vector<std::string>
Default value: []
Flat list of hosts of one unsharded table, all sharing port. For exactly-once, this form is used with TClickHouseBatchingSink. Required unless host or shard_hosts is set. Every host must be reachable and match the engine and complete ordered schema (name, type, default_kind, default_expression). For Replicated*MergeTree, zookeeper_name and zookeeper_path must match, while replica_name may differ. Multiple hosts with plain MergeTree or SharedMergeTree are rejected; single-host SharedMergeTree remains supported. This target-identity rule applies to all delivery guarantees. A single-entry list is rejected — use host instead.
|
|
host_selection_policy
|
Type: NYT::NFlow::EClickHouseHostSelectionPolicy
Default value: ordered_round_robin
Endpoint order within each shard. The default ordered_round_robin preserves the configured order. random_start chooses a uniformly distributed starting position once when a client is constructed and rotates the list. clickhouse-cpp then exhausts all endpoints in round-robin order. The policy does not change row routing, deduplication tokens, or the topology fingerprint.
|
|
shard_hosts
|
Type: THashMap<std::string, std::vector<std::string>>
Default value: {}
Mapping from a shard name to the list of that shard's hosts, used by the multi-shard form. Required by TShardedClickHouseBatchingSink and rejected by TClickHouseBatchingSink; the other sink classes accept any of the three host forms. Every endpoint must be reachable. Within each shard, the engine and complete ordered schema (name, type, default_kind, default_expression) must match; schemas must also match across shards. For replicas, zookeeper_name and zookeeper_path must match, while replica_name may differ. Multiple hosts with plain MergeTree or SharedMergeTree are rejected; single-host SharedMergeTree remains supported. This target-identity rule applies to all delivery guarantees. Shard names must be non-empty and consist of letters, digits, _ and -; a shard's host list must not be empty. A shard name is a permanent identifier: it is part of the deduplication token.
|
|
sharding_key_columns
|
Type: std::vector<std::string>
Default value: []
Columns whose values produce the routing key. Requires shard_hosts: the unsharded forms reject this parameter during spec validation. An empty list routes by message id, that is, without co-location.
|
|
user
|
Type: std::string
Default value: default
ClickHouse user.
|
|
password_env_var
|
Type: std::string
Default value: TString("")
Name of the environment variable holding the password; resolved via GetEnv when the client is created. An empty value means connecting without a password.
|
|
database
|
Type: std::string
Default value: default
Database of the target table. Shared by all shards.
|
|
table
|
Type: std::string
Required parameter
Target table. Shared by all shards, so <database>.<table> must exist on every host of every shard.
|
|
codec
|
Type: NYT::NFlow::EClickHouseCodec
Default value: lz4
Native protocol compression.
|
|
enable_tls
|
Type: bool
Default value: false
Connect over TCP+TLS instead of plain TCP. The port is not changed implicitly: set port to the server's secure native port (usually 9440).
|
|
tls_ca_files
|
Type: std::vector<std::string>
Default value: []
Paths to root certificate (CA) files used to verify the server certificate. An empty list is valid when TLS is disabled; when TLS is enabled, it means that the system root certificates are used. A nonempty list requires enable_tls = true; otherwise configuration validation fails.
|
|
tls_ca_directory
|
Type: std::string
Default value: TString("")
Path to a directory of root certificates (CA) used to verify the server certificate. An empty value is valid when TLS is disabled. A nonempty path requires enable_tls = true; otherwise configuration validation fails.
|
|
tls_skip_verification
|
Type: bool
Default value: false
Skip TLS session verification (the server certificate and so on). Insecure; use only for test environments with self-signed certificates. Setting this option to true requires enable_tls = true; otherwise configuration validation fails. The default value, false, is valid when TLS is disabled.
|
NYT::NFlow::TUnitedParameters<NYT::NFlow::TPassthroughComputation>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
NYT::NFlow::TUnitedParameters<NYT::NFlow::TQueueSource>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
finite
|
Type: bool
Default value: false
Treat the source as finite: remember the number of messages in it on startup and transition the stream to the completed state after all those messages are read.
|
|
queue_path
|
Type: NYT::NYPath::TRichYPath
Required parameter
Path to the queue including the cluster.
|
|
consumer_path
|
Type: NYT::NYPath::TRichYPath
Required parameter
Path to the queue consumer with the cluster specified.
|
|
try_parse_flow_queue_meta
|
Type: bool
Default value: true
Whether to retrieve and parse meta from the queue (the meta contains watermarks provided by the writer to the queue).
|
|
flow_queue_meta_column
|
Type: std::string
Default value: flow_queue_meta
Column from which to retrieve meta.
|
|
ignore_malformed_flow_queue_meta
|
Type: bool
Default value: false
Whether to ignore invalid meta or fail.
|
|
partition_filter
|
Type: std::optional<std::vector<std::pair<int, int>>>
|
Additional parameters
|
update_info_period
|
Type: TDuration
Default value: 15s
Period for refreshing the source's auxiliary partition information. A connector may use this tick for status requests and read-session liveness checks.
|
|
byte_size_alpha
|
Type: double
Default value: 0.05
Exponential smoothing coefficient for the average byte total and message count per offset: larger values make the estimate react to new data faster.
|
|
update_partition_count_period
|
Type: TDuration
Default value: 1m
How often the controller should update the partition count in the queue (the controller changes the computation partition count according to the queue partition count).
|
|
update_partition_count_retry_min_backoff
|
Type: TDuration
Default value: 1s
Initial delay before retrying a failed partition count update, subject to jitter. The value is capped by update_partition_count_period; subsequent delays increase exponentially.
|
NYT::NFlow::TUnitedParameters<NYT::NFlow::TRandomSource>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
finite
|
Type: bool
Default value: false
Treat the source as finite: remember the number of messages in it on startup and transition the stream to the completed state after all those messages are read.
|
Additional parameters
|
update_info_period
|
Type: TDuration
Default value: 15s
Period for refreshing the source's auxiliary partition information. A connector may use this tick for status requests and read-session liveness checks.
|
|
byte_size_alpha
|
Type: double
Default value: 0.05
Exponential smoothing coefficient for the average byte total and message count per offset: larger values make the estimate react to new data faster.
|
NYT::NFlow::TUnitedParameters<NYT::NFlow::TServiceLogSource>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
finite
|
Type: bool
Default value: false
Treat the source as finite: remember the number of messages in it on startup and transition the stream to the completed state after all those messages are read.
|
|
table_joiner
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TTableJoinerSpec>
Default value: {}
Joiner spec.
|
Additional parameters
|
update_info_period
|
Type: TDuration
Default value: 15s
Period for refreshing the source's auxiliary partition information. A connector may use this tick for status requests and read-session liveness checks.
|
|
byte_size_alpha
|
Type: double
Default value: 0.05
Exponential smoothing coefficient for the average byte total and message count per offset: larger values make the estimate react to new data faster.
|
NYT::NFlow::TUnitedParameters<NYT::NFlow::TShardedClickHouseBatchingSink>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
host
|
Type: std::string
Default value: TString("")
The single ClickHouse host of the unsharded form. For exactly-once, this form is used with TClickHouseBatchingSink. Required unless hosts or shard_hosts is set: exactly one of the three host forms must be present in the spec. The host must be reachable and expose a table with a supported engine and the complete ordered schema (name, type, default_kind, default_expression).
|
|
port
|
Type: unsigned short
Default value: 9000
Native TCP port. Shared by all shards.
|
|
hosts
|
Type: std::vector<std::string>
Default value: []
Flat list of hosts of one unsharded table, all sharing port. For exactly-once, this form is used with TClickHouseBatchingSink. Required unless host or shard_hosts is set. Every host must be reachable and match the engine and complete ordered schema (name, type, default_kind, default_expression). For Replicated*MergeTree, zookeeper_name and zookeeper_path must match, while replica_name may differ. Multiple hosts with plain MergeTree or SharedMergeTree are rejected; single-host SharedMergeTree remains supported. This target-identity rule applies to all delivery guarantees. A single-entry list is rejected — use host instead.
|
|
host_selection_policy
|
Type: NYT::NFlow::EClickHouseHostSelectionPolicy
Default value: ordered_round_robin
Endpoint order within each shard. The default ordered_round_robin preserves the configured order. random_start chooses a uniformly distributed starting position once when a client is constructed and rotates the list. clickhouse-cpp then exhausts all endpoints in round-robin order. The policy does not change row routing, deduplication tokens, or the topology fingerprint.
|
|
shard_hosts
|
Type: THashMap<std::string, std::vector<std::string>>
Default value: {}
Mapping from a shard name to the list of that shard's hosts, used by the multi-shard form. Required by TShardedClickHouseBatchingSink and rejected by TClickHouseBatchingSink; the other sink classes accept any of the three host forms. Every endpoint must be reachable. Within each shard, the engine and complete ordered schema (name, type, default_kind, default_expression) must match; schemas must also match across shards. For replicas, zookeeper_name and zookeeper_path must match, while replica_name may differ. Multiple hosts with plain MergeTree or SharedMergeTree are rejected; single-host SharedMergeTree remains supported. This target-identity rule applies to all delivery guarantees. Shard names must be non-empty and consist of letters, digits, _ and -; a shard's host list must not be empty. A shard name is a permanent identifier: it is part of the deduplication token.
|
|
sharding_key_columns
|
Type: std::vector<std::string>
Default value: []
Columns whose values produce the routing key. Requires shard_hosts: the unsharded forms reject this parameter during spec validation. An empty list routes by message id, that is, without co-location.
|
|
user
|
Type: std::string
Default value: default
ClickHouse user.
|
|
password_env_var
|
Type: std::string
Default value: TString("")
Name of the environment variable holding the password; resolved via GetEnv when the client is created. An empty value means connecting without a password.
|
|
database
|
Type: std::string
Default value: default
Database of the target table. Shared by all shards.
|
|
table
|
Type: std::string
Required parameter
Target table. Shared by all shards, so <database>.<table> must exist on every host of every shard.
|
|
codec
|
Type: NYT::NFlow::EClickHouseCodec
Default value: lz4
Native protocol compression.
|
|
enable_tls
|
Type: bool
Default value: false
Connect over TCP+TLS instead of plain TCP. The port is not changed implicitly: set port to the server's secure native port (usually 9440).
|
|
tls_ca_files
|
Type: std::vector<std::string>
Default value: []
Paths to root certificate (CA) files used to verify the server certificate. An empty list is valid when TLS is disabled; when TLS is enabled, it means that the system root certificates are used. A nonempty list requires enable_tls = true; otherwise configuration validation fails.
|
|
tls_ca_directory
|
Type: std::string
Default value: TString("")
Path to a directory of root certificates (CA) used to verify the server certificate. An empty value is valid when TLS is disabled. A nonempty path requires enable_tls = true; otherwise configuration validation fails.
|
|
tls_skip_verification
|
Type: bool
Default value: false
Skip TLS session verification (the server certificate and so on). Insecure; use only for test environments with self-signed certificates. Setting this option to true requires enable_tls = true; otherwise configuration validation fails. The default value, false, is valid when TLS is disabled.
|
NYT::NFlow::TUnitedParameters<NYT::NFlow::TSwiftPassthroughComputation>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
allow_batching_with_relaxed_guarantees
|
Type: bool
Default value: false
Allows combining several input messages into one output message (batching). Useful for folding many small messages into one larger message to reduce the message count load on downstream partitions. Default is %false: each output message is produced by exactly one parent. MessageId is deterministically inherited from the parent, providing exactly-once semantics with deterministic user functions. When set to %true, an output message can have multiple parents. A parent is considered processed only when all its children are processed. On restarts, batch boundaries may differ, so the same parent message could end up in multiple different batched children — downstream computations must be prepared to see each parent's contents more than once (at-least-once semantics). Additionally, batched output messages receive a MessageId derived deterministically from the set of parent MessageIds: a replay of the same batch yields the same MessageId and is deduplicated, while a batch with a different composition yields a new one. This means that the order of messages by MessageId within a single key on the downstream Computation side can be disrupted. If subsequent pipeline stages rely on MessageId ordering within a key, you must rewrite that logic to account for this.
|
NYT::NFlow::TUnitedParameters<NYT::NFlow::TSwiftPassthroughOrderedSourceComputation>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
This structure has no main parameters.
NYT::NFlow::TUnitedParameters<NYT::NFlow::TSyncQueueSink>
Source: yt/yt/flow/library/cpp/common/registry-inl.h
|
Parameter
|
Description
|
|
queue_path
|
Type: NYT::NYPath::TRichYPath
Required parameter
Path to the queue including the cluster.
|
|
write_flow_queue_meta
|
Type: bool
Default value: false
Whether to write meta to the queue (watermarks are written in meta).
|
|
flow_queue_meta_column
|
Type: std::string
Default value: flow_queue_meta
Column to which to write meta.
|
|
tablet_index_expression
|
Type: std::optional<std::string>
Verbatim routing mode. A QL expression (builtins only, e.g. farm_hash) over the message columns whose value is written to the system column $tablet_index verbatim. Mutually exclusive with tablet_index_routing_hash_expression; when neither expression is set, routing is off ($tablet_index is not written and the tablet is chosen by the driver). Sharded queues are inconvenient to operate (object→tablet affinity is tightly coupled and a reshard remaps keys), so enable routing only when genuinely needed. Routing is supported only on sync queue sinks; async sinks reject the routing params at spec load (YTFLOW-766).
|
|
tablet_index_routing_hash_expression
|
Type: std::optional<std::string>
Hash routing mode. A QL expression (builtins only, e.g. farm_hash) yielding a uint64 hash that is reduced to a tablet index by tablet_index_routing_hash_policy over tablet_count. Mutually exclusive with tablet_index_expression. Supported only on sync queue sinks; async sinks reject the routing params (YTFLOW-766).
|
|
tablet_index_routing_hash_policy
|
Type: std::optional<NYT::NFlow::EQueueTabletIndexRoutingHashPolicy>
The policy for reducing the tablet_index_routing_hash_expression hash to a tablet index. Required together with it. range — contiguous equal-width hash ranges (rangeSize = 2^64 / tablet_count), recommended: a consumer range-partitioned by the same key reads only its own tablet. modulo — hash % tablet_count; discouraged: Flow computations are range-partitioned by key, so a modulo-partitioned queue forces every reader to read every tablet (a full mesh on read).
|
|
tablet_count
|
Type: std::optional<long>
The number of tablets used to reduce tablet_index_routing_hash_expression to a tablet index. Optional. When unset, the sink periodically re-reads the target queue's @tablet_count (from the master cache) and follows a reshard without a restart. When set explicitly, the value is fixed (no queue read) and a reshard is reflected only by changing the config.
|
|
column_filter
|
Type: std::optional<THashSet<std::string>>
Which columns from the message to write to the queue (default — all).
|
Additional parameters
|
update_partition_count_period
|
Type: TDuration
Default value: 1m
How often the controller should update the partition count in the queue (the controller changes the computation partition count according to the queue partition count).
|
|
update_partition_count_retry_min_backoff
|
Type: TDuration
Default value: 1s
Initial delay before retrying a failed partition count update, subject to jitter. The value is capped by update_partition_count_period; subsequent delays increase exponentially.
|
NYT::NFlow::TVanillaConfig
Source: yt/yt/flow/library/cpp/runner/vanilla_launcher.h
|
Parameter
|
Description
|
|
enable
|
Type: bool
Default value: false
Enable vanilla mode. If %false, the block is ignored.
|
|
pool
|
Type: std::string
Required parameter
The pool in which the vanilla operation runs.
|
|
worker
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TVanillaTaskConfig>
Required parameter
Config for worker jobs (number of jobs and resource limits). You must set count.
|
|
controller
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TVanillaTaskConfig>
Default value: {'count': 1}
Config for controller jobs (number of jobs and resource limits). If not set, one job with default limits is used.
|
|
runtime_cluster
|
Type: std::optional<std::string>
The cluster on which the vanilla operation runs. Defaults to the pipeline's cluster_url.
|
|
runtime_proxy_role
|
Type: std::optional<std::string>
The RPC proxy role for the runtime_cluster (the pipeline cluster's role may not exist on that cluster). Taken into account only when runtime_cluster differs from the pipeline cluster.
|
|
cache_path
|
Type: TString
Default value: //tmp/yt_wrapper/file_storage/new_cache
A content-addressed cache to which job files are uploaded (shared by all flow operations on the cluster). A long-term copy in the pipeline folder is a cheap CopyNode from this cache.
|
|
max_failed_job_count
|
Type: unsigned long
Default value: 10000
The maximum number of failed jobs after which the vanilla operation fails.
|
|
max_stderr_count
|
Type: int
Default value: 150
|
|
wait_timeout
|
Type: TDuration
Default value: 5m
Timeout for waiting for pipeline states during a graceful stop of the previous vanilla operation.
|
|
alias
|
Type: std::optional<std::string>
An explicit alias for the vanilla operation. If not set, it is generated as *flow-runner <cluster>:<path>.
|
|
title
|
Type: std::optional<std::string>
Title of the vanilla operation.
|
|
network_project
|
Type: std::optional<std::string>
The network project for a vanilla operation has no default in open-source builds.
|
|
proxy_url_aliasing_rules
|
Type: THashMap<std::string, std::string>
Default value: {}
Aliases for proxy URLs, injected into the flow-server inside the vanilla jobs.
|
|
secret_env
|
Type: std::vector<std::string>
Default value: []
|
|
node_config
|
Type: NYT::TIntrusivePtr<NYT::NYTree::INode>
An arbitrary YSON patch applied on top of the default TFlowNodeConfig inside the vanilla jobs.
|
NYT::NFlow::TVanillaTaskConfig
Source: yt/yt/flow/library/cpp/runner/vanilla_launcher.h
|
Parameter
|
Description
|
|
count
|
Type: int
Required parameter
Number of jobs for the task (controller or worker) in the vanilla operation.
|
|
memory_limit
|
Type: std::optional<NYT::NYTree::TSize>
Memory limit per job. If not set, a single shared default of 18 GiB is used, the same for the controller and worker.
|
|
cpu_limit
|
Type: std::optional<int>
CPU limit per job. If not set, a single shared default of 6 is used, the same for the controller and worker.
|
|
set_container_cpu_limit
|
Type: bool
Default value: false
Request a container CPU ceiling equal to cpu_limit on execution backends that support it. When %false, the runner leaves the operation field unset.
|
|
port_count
|
Type: std::optional<int>
Number of ports to request from YT (ports are provided via YT_PORT_<i> and override fixed ports from node-config): YT_PORT_0 is rpc_port (and bus_server.port), YT_PORT_1 is monitoring_port, YT_PORT_2 is companion.port. Needed on hosts with a shared network where fixed ports of neighboring jobs would collide: 2 for the controller and a worker without a companion, 3 for a worker with a companion (Python, Java, Go, C++ companion). A launch without a network project fills it in automatically (2 for the controller, 3 for the worker); an explicit 0 keeps the fixed ports. Once the worker's port_count is positive, the fixed companion.port = 10082 is no longer filled in: a worker with a companion and fewer than three ports fails to start instead of dialing a neighboring job's companion.
|
|
local_files
|
Type: THashMap<std::string, std::string>
Default value: {}
|
|
cypress_files
|
Type: THashMap<std::string, std::string>
Default value: {}
|
|
layers
|
Type: std::vector<std::string>
Default value: []
Cypress paths of the porto layers mounted into the task's root filesystem. A non-empty list on at least one task enables porto jobs for the whole vanilla operation.
|
|
system_layer_path
|
Type: std::optional<std::string>
The task's base OS layer; overrides the default system layer.
|
|
docker_image
|
Type: std::optional<std::string>
Docker image for the task's root filesystem, for clusters whose job environment pulls images rather than mounting porto layers. Mutually exclusive with layers.
|
NYT::NFlow::TWatermarkAlignmentSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
group_name
|
Type: NYT::TStrongTypedef<std::string, NYT::NFlow::TWatermarkAlignmentGroupTag, NYT::TStrongTypedefOptions{true}>
Default value: default-alignment-group
The group name.
|
|
drift_bound
|
Type: TDuration
Default value: 20m
The maximum possible lead of an individual partition's watermark over the group's watermark.
|
|
read_delays
|
Type: std::optional<THashMap<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TStreamIdTag>, TDuration>>
An advanced use option. If the EventTimestamp of an output message is greater than EventWatermark - Delay for at least one specified stream, reading is paused. If misconfigured, this option can cause a complete read stall.
|
NYT::NFlow::TWatermarkGeneratorSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
use_source_watermark
|
Type: bool
Default value: false
Use SourceWatermark for EventWatermark generation. This option completely disables all heuristics for determining EventWatermark. Use it strictly when the input stream contains special markers with EventWatermark.
|
|
out_of_orderness_bound
|
Type: TDuration
Default value: 0
The upper bound on possible event reordering, used for estimating EventWatermark. Event reordering is not accounted for by default.
|
|
idle_partitions
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TIdlePartitionsSpec>
Heuristic settings for ignoring partitions with zero write throughput.
|
|
late_data_partitions
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TLateDataPartitionsSpec>
|
NYT::NFlow::TWatermarkPercentileSpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
value
|
Type: NYT::TStrongTypedef<double, NYT::NFlow::TWatermarkPercentileTag, NYT::TStrongTypedefOptions{true}>
Default value: 100.0
The percentage of events unconditionally considered for EventWatermark calculation.
|
|
delay
|
Type: TDuration
Default value: 1m
The freshness threshold for events below the selected percentile.
|
NYT::NFlow::TWatermarkStrategySpec
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter
|
Description
|
|
event_timestamp_assigner
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TEventTimestampAssignerSpec>
Default value: {}
Automatic population of EventTimestamp for all output messages based on the value of a specified column.
|
|
watermark_generator
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TWatermarkGeneratorSpec>
Settings for the EventWatermark generation algorithm. For a Computation with at least one source, omitting this parameter is equivalent to an empty block with default settings. A multi-source Computation is processed as independent single-source computations: for each source, unavailable-group, idle-partition, and late-data-partition heuristics run independently in that order, and the prepared partitions are merged only afterwards. idle_partitions.max_ratio is applied separately to each source, and availability-group suppression decisions in one source do not affect another. Therefore, a fully idle source is not treated as a partial-idle stall and does not by itself raise /idle_partitions_watermark_stall, even if it gates the merged Computation watermark while another source remains active. Availability-group identifiers do not need to match across sources.
|
|
watermark_alignment
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TWatermarkAlignmentSpec>
Settings for aligning source across different partitions, including those of different Computation instances.
|
|
watermark_percentile
|
Type: NYT::TIntrusivePtr<NYT::NFlow::TWatermarkPercentileSpec>
Default value: {}
Settings for a "percentile" EventWatermark. Used to ignore the oldest events within a configured percentage.
|
NYT::NHttp::TClientConfig
Source: yt/yt/core/http/config.h
|
Parameter
|
Description
|
|
read_buffer_size
|
Type: int
Default value: 131072
|
|
max_redirect_count
|
Type: int
Default value: 0
|
|
ignore_continue_responses
|
Type: bool
Default value: false
|
|
connection_idle_timeout
|
Type: TDuration
Default value: 5m
|
|
header_read_timeout
|
Type: TDuration
Default value: 30s
|
|
body_read_idle_timeout
|
Type: TDuration
Default value: 5m
|
|
write_idle_timeout
|
Type: TDuration
Default value: 5m
|
|
max_idle_connections
|
Type: int
Default value: 0
|
|
dns_resolve_options
|
Type: std::optional<NYT::NDns::TDnsResolveOptions>
|
|
omit_question_mark_for_empty_query
|
Type: bool
Default value: false
|
|
dialer
|
Type: NYT::TIntrusivePtr<NYT::NNet::TDialerConfig>
Default value: {}
|
NYT::NHttps::TClientConfig
Source: yt/yt/core/https/config.h
|
Parameter
|
Description
|
|
read_buffer_size
|
Type: int
Default value: 131072
|
|
max_redirect_count
|
Type: int
Default value: 0
|
|
ignore_continue_responses
|
Type: bool
Default value: false
|
|
connection_idle_timeout
|
Type: TDuration
Default value: 5m
|
|
header_read_timeout
|
Type: TDuration
Default value: 30s
|
|
body_read_idle_timeout
|
Type: TDuration
Default value: 5m
|
|
write_idle_timeout
|
Type: TDuration
Default value: 5m
|
|
max_idle_connections
|
Type: int
Default value: 0
|
|
dns_resolve_options
|
Type: std::optional<NYT::NDns::TDnsResolveOptions>
|
|
omit_question_mark_for_empty_query
|
Type: bool
Default value: false
|
|
dialer
|
Type: NYT::TIntrusivePtr<NYT::NNet::TDialerConfig>
Default value: {}
|
|
credentials
|
Type: NYT::TIntrusivePtr<NYT::NHttps::TClientCredentialsConfig>
|
|
allow_http
|
Type: bool
Default value: false
|
NYT::NHttps::TClientCredentialsConfig
Source: yt/yt/core/https/config.h
NYT::NLogging::ELogFamily
|
Possible values
|
Description
|
|
plain_text
|
|
|
structured
|
|
NYT::NLogging::ELogLevel
|
Possible values
|
Description
|
|
minimum
|
|
|
trace
|
|
|
debug
|
|
|
info
|
|
|
warning
|
|
|
error
|
|
|
alert
|
|
|
fatal
|
|
|
maximum
|
|
NYT::NLogging::TLogManagerConfig
Source: yt/yt/core/logging/config.h
|
Parameter
|
Description
|
|
flush_period
|
Type: std::optional<TDuration>
|
|
watch_period
|
Type: std::optional<TDuration>
|
|
check_space_period
|
Type: std::optional<TDuration>
|
|
rotation_check_period
|
Type: TDuration
Default value: 5s
|
|
min_disk_space
|
Type: long
Default value: 5368709120
|
|
high_backlog_watermark
|
Type: int
Default value: 10000000
|
|
low_backlog_watermark
|
Type: int
Default value: 1000000
|
|
shutdown_grace_timeout
|
Type: TDuration
Default value: 1s
|
|
shutdown_busy_timeout
|
Type: TDuration
Default value: 0
|
|
rules
|
Type: std::vector<NYT::TIntrusivePtr<NYT::NLogging::TRuleConfig>>
Required parameter
|
|
writers
|
Type: THashMap<std::string, NYT::TIntrusivePtr<NYT::NYTree::IMapNode>>
Required parameter
|
|
category_rate_limits
|
Type: THashMap<std::string, long>
Default value: {}
|
|
suppressed_messages
|
Type: std::vector<std::string>
Default value: []
|
|
message_level_overrides
|
Type: THashMap<std::string, NYT::NLogging::ELogLevel>
Default value: {}
|
|
request_suppression_timeout
|
Type: TDuration
Default value: 0
|
|
enable_anchor_profiling
|
Type: bool
Default value: false
|
|
min_logged_message_rate_to_profile
|
Type: double
Default value: 1.0
|
|
abort_on_alert
|
Type: bool
Default value: false
|
|
structured_validation_sampling_rate
|
Type: double
Default value: 0.01
|
|
compression_thread_count
|
Type: int
Default value: 1
|
NYT::NLogging::TLogManagerDynamicConfig
Source: yt/yt/core/logging/config.h
|
Parameter
|
Description
|
|
min_disk_space
|
Type: std::optional<long>
|
|
high_backlog_watermark
|
Type: std::optional<int>
|
|
low_backlog_watermark
|
Type: std::optional<int>
|
|
rules
|
Type: std::optional<std::vector<NYT::TIntrusivePtr<NYT::NLogging::TRuleConfig>>>
|
|
category_rate_limits
|
Type: std::optional<THashMap<std::string, long>>
|
|
suppressed_messages
|
Type: std::optional<std::vector<std::string>>
|
|
message_level_overrides
|
Type: THashMap<std::string, NYT::NLogging::ELogLevel>
Default value: {}
|
|
request_suppression_timeout
|
Type: std::optional<TDuration>
|
|
enable_anchor_profiling
|
Type: std::optional<bool>
|
|
min_logged_message_rate_to_profile
|
Type: std::optional<double>
|
|
abort_on_alert
|
Type: std::optional<bool>
|
|
structured_validation_sampling_rate
|
Type: std::optional<double>
|
|
compression_thread_count
|
Type: std::optional<int>
|
NYT::NLogging::TRuleConfig
Source: yt/yt/core/logging/config.h
|
Parameter
|
Description
|
|
include_categories
|
Type: std::optional<THashSet<std::string>>
|
|
exclude_categories
|
Type: THashSet<std::string>
Default value: []
|
|
min_level
|
Type: NYT::NLogging::ELogLevel
Default value: minimum
|
|
max_level
|
Type: NYT::NLogging::ELogLevel
Default value: maximum
|
|
family
|
Type: std::optional<NYT::NLogging::ELogFamily>
|
|
writers
|
Type: std::vector<std::string>
Required parameter
|
NYT::NNet::TAddressResolverConfig
Source: yt/yt/core/net/config.h
|
Parameter
|
Description
|
|
shard_count
|
Type: unsigned long
Default value: 1
|
|
expire_after_access_time
|
Type: TDuration
Default value: 5m
|
|
expire_after_successful_update_time
|
Type: TDuration
Default value: 2m
|
|
expire_after_failed_update_time
|
Type: TDuration
Default value: 30s
|
|
refresh_time
|
Type: std::optional<TDuration>
Default value: 60000
|
|
expiration_period
|
Type: std::optional<TDuration>
Default value: 60000
|
|
batch_update
|
Type: bool
Default value: false
|
|
retries
|
Type: int
Default value: 25
|
|
retry_delay
|
Type: TDuration
Default value: 200ms
|
|
resolve_timeout
|
Type: TDuration
Default value: 1s
|
|
max_resolve_timeout
|
Type: TDuration
Default value: 15s
|
|
warning_timeout
|
Type: TDuration
Default value: 3s
|
|
jitter
|
Type: std::optional<double>
Default value: 0.5
|
|
force_tcp
|
Type: bool
Default value: false
|
|
keep_socket
|
Type: bool
Default value: true
|
|
enable_ipv4
|
Type: bool
Default value: false
|
|
enable_ipv6
|
Type: bool
Default value: true
|
|
localhost_name_override
|
Type: std::optional<std::string>
|
|
resolve_hostname_into_fqdn
|
Type: bool
Default value: true
|
NYT::NNet::TDialerConfig
Source: yt/yt/core/net/config.h
|
Parameter
|
Description
|
|
enable_no_delay
|
Type: bool
Default value: true
|
|
enable_aggressive_reconnect
|
Type: bool
Default value: false
|
|
allow_bypass_tls
|
Type: bool
Default value: false
|
|
min_rto
|
Type: TDuration
Default value: 100ms
|
|
max_rto
|
Type: TDuration
Default value: 30s
|
|
rto_scale
|
Type: double
Default value: 2.0
|
|
connect_timeout
|
Type: TDuration
Default value: 15s
|
NYT::NPipeIO::TPipeIODispatcherConfig
Source: yt/yt/library/pipe_io/config.h
|
Parameter
|
Description
|
|
thread_pool_polling_period
|
Type: TDuration
Default value: 10ms
|
NYT::NPipeIO::TPipeIODispatcherDynamicConfig
Source: yt/yt/library/pipe_io/config.h
|
Parameter
|
Description
|
|
thread_pool_polling_period
|
Type: std::optional<TDuration>
|
NYT::NProfiling::ELabelSanitizationPolicy
|
Possible values
|
Description
|
|
none
|
|
|
weak
|
|
|
strong
|
|
NYT::NProfiling::TResourceTrackerConfig
Source: yt/yt/library/profiling/resource_tracker/config.h
|
Parameter
|
Description
|
|
enable
|
Type: bool
Default value: true
|
|
cpu_to_vcpu_factor
|
Type: std::optional<double>
|
NYT::NProfiling::TShardConfig
Source: yt/yt/library/profiling/solomon/config.h
|
Parameter
|
Description
|
|
filter
|
Type: std::vector<std::string>
Default value: []
|
|
grid_step
|
Type: std::optional<TDuration>
|
|
strip_sensors_name_prefix
|
Type: bool
Default value: false
|
NYT::NProfiling::TSolomonExporterConfig
Source: yt/yt/library/profiling/solomon/config.h
|
Parameter
|
Description
|
|
enable
|
Type: bool
Default value: true
|
|
grid_step
|
Type: TDuration
Default value: 5s
|
|
linger_timeout
|
Type: TDuration
Default value: 5m
|
|
window_size
|
Type: int
Default value: 12
|
|
thread_pool_size
|
Type: int
Default value: 1
|
|
encoding_thread_pool_size
|
Type: int
Default value: 1
|
|
thread_pool_polling_period
|
Type: TDuration
Default value: 10ms
|
|
encoding_thread_pool_polling_period
|
Type: TDuration
Default value: 10ms
|
|
convert_counters_to_rate_for_solomon
|
Type: bool
Default value: true
|
|
rename_converted_counters
|
Type: bool
Default value: true
|
|
convert_counters_to_delta_gauge
|
Type: bool
Default value: false
|
|
enable_histogram_compat
|
Type: bool
Default value: false
|
|
split_rate_histogram_into_gauges
|
Type: bool
Default value: false
|
|
report_timestamps_for_rate_metrics
|
Type: bool
Default value: true
|
|
export_summary
|
Type: bool
Default value: false
|
|
export_summary_as_sum
|
Type: bool
Default value: false
|
|
export_summary_as_max
|
Type: bool
Default value: true
|
|
export_summary_as_min
|
Type: bool
Default value: false
|
|
export_summary_as_avg
|
Type: bool
Default value: false
|
|
mark_aggregates
|
Type: bool
Default value: true
|
|
enable_solomon_aggregates
|
Type: bool
Default value: false
|
|
export_globals_as_mem_only
|
Type: bool
Default value: false
|
|
strip_sensors_name_prefix
|
Type: bool
Default value: false
|
|
enable_self_profiling
|
Type: bool
Default value: true
|
|
report_build_info
|
Type: bool
Default value: true
|
|
report_kernel_version
|
Type: bool
Default value: true
|
|
report_restart
|
Type: bool
Default value: true
|
|
read_delay
|
Type: TDuration
Default value: 5s
|
|
host
|
Type: std::optional<std::string>
|
|
instance_tags
|
Type: THashMap<std::string, std::string>
Default value: {}
|
|
shards
|
Type: THashMap<std::string, NYT::TIntrusivePtr<NYT::NProfiling::TShardConfig>>
Default value: {}
|
|
response_cache_ttl
|
Type: TDuration
Default value: 2m
|
|
update_sensor_service_tree_period
|
Type: TDuration
Default value: 30s
|
|
producer_collection_batch_size
|
Type: int
Default value: 100
|
|
label_sanitization_policy
|
Type: NYT::NProfiling::ELabelSanitizationPolicy
Default value: none
|
NYT::NProfiling::TSolomonProxyConfig
Source: yt/yt/library/profiling/solomon/proxy.h
|
Parameter
|
Description
|
|
public_component_names
|
Type: std::optional<THashSet<std::string>>
|
|
max_endpoints_per_request
|
Type: int
Default value: 30
|
NYT::NProfiling::TSolomonRegistryConfig
Source: yt/yt/library/profiling/solomon/config.h
|
Parameter
|
Description
|
|
enable_rseq
|
Type: bool
Default value: false
|
NYT::NProfiling::TSolomonRegistryDynamicConfig
Source: yt/yt/library/profiling/solomon/config.h
|
Parameter
|
Description
|
|
enable_rseq
|
Type: std::optional<bool>
|
NYT::NQueryClient::EStatisticsAggregation
|
Possible values
|
Description
|
|
none
|
|
|
depth
|
|
|
depth_omit_node
|
|
NYT::NQueryClient::TCodegenCacheConfig
Source: yt/yt/library/query/engine_api/cg_cache.h
|
Parameter
|
Description
|
|
capacity
|
Type: long
Default value: 512
|
|
younger_size_fraction
|
Type: double
Default value: 0.25
|
|
shard_count
|
Type: int
Default value: 1
|
|
touch_buffer_capacity
|
Type: int
Default value: 65536
|
|
small_ghost_cache_ratio
|
Type: double
Default value: 0.5
|
|
large_ghost_cache_ratio
|
Type: double
Default value: 2.0
|
|
enable_ghost_caches
|
Type: bool
Default value: true
|
|
reject_oversized_items
|
Type: bool
Default value: false
|
NYT::NQueryClient::TCodegenCacheDynamicConfig
Source: yt/yt/library/query/engine_api/cg_cache.h
|
Parameter
|
Description
|
|
capacity
|
Type: std::optional<long>
|
|
younger_size_fraction
|
Type: std::optional<double>
|
|
enable_ghost_caches
|
Type: bool
Default value: true
|
|
reject_oversized_items
|
Type: std::optional<bool>
|
NYT::NQueryClient::TQueryEngineConfig
Source: yt/yt/library/query/engine_api/query_engine_config.h
NYT::NQueryClient::TQueryEngineDynamicConfig
Source: yt/yt/library/query/engine_api/query_engine_config.h
|
Parameter
|
Description
|
|
codegen_cache
|
Type: NYT::TIntrusivePtr<NYT::NQueryClient::TCodegenCacheDynamicConfig>
Default value: {}
|
|
statistics_aggregation
|
Type: std::optional<NYT::NQueryClient::EStatisticsAggregation>
|
|
use_order_by_in_join_subqueries
|
Type: std::optional<bool>
|
|
enable_parallelize_unordered_group_by
|
Type: std::optional<bool>
|
|
expression_builder_version
|
Type: std::optional<int>
|
|
max_projection_count
|
Type: int
Default value: 1024
|
|
codegen_optimization_level
|
Type: std::optional<NYT::NCodegen::EOptimizationLevel>
|
|
allow_udf_object_code_cache
|
Type: std::optional<bool>
|
|
rewrite_cardinality_into_hyper_log_log_with_precision
|
Type: std::optional<bool>
|
|
allow_join_with_async_last_committed_timestamp_if_require_sync_replica_is_false
|
Type: std::optional<bool>
|
|
truncated_query_length_for_tracing
|
Type: std::optional<int>
|
|
allow_reverse_scan_for_order_by
|
Type: std::optional<bool>
|
|
allow_heavy_range_inference_in_joins
|
Type: std::optional<bool>
|
|
prefetch_join_tables
|
Type: std::optional<bool>
|
|
allow_multiple_join_subqueries_for_non_lookup_joins
|
Type: std::optional<bool>
|
|
join_cache_size
|
Type: std::optional<long>
|
|
max_subsplits_per_tablet
|
Type: std::optional<int>
|
NYT::NRpc::EPeerPriorityStrategy
|
Possible values
|
Description
|
|
none
|
|
|
prefer_local
|
|
NYT::NRpc::ERequestTracingMode
|
Possible values
|
Description
|
|
enable
|
|
|
disable
|
|
|
force
|
|
NYT::NRpc::NGrpc::TChannelConfig
Source: yt/yt/core/rpc/grpc/config.h
|
Parameter
|
Description
|
|
credentials
|
Type: NYT::TIntrusivePtr<NYT::NRpc::NGrpc::TChannelCredentialsConfig>
|
|
grpc_arguments
|
Type: THashMap<std::string, NYT::TIntrusivePtr<NYT::NYTree::INode>>
Default value: {}
|
|
address
|
Type: std::string
Default value: TString("")
|
NYT::NRpc::NGrpc::TChannelCredentialsConfig
Source: yt/yt/core/rpc/grpc/config.h
NYT::NRpc::NGrpc::TDispatcherConfig
Source: yt/yt/core/rpc/grpc/config.h
|
Parameter
|
Description
|
|
dispatcher_thread_count
|
Type: int
Default value: 4
|
|
grpc_thread_count
|
Type: int
Default value: 4
|
|
grpc_event_engine_thread_count
|
Type: int
Default value: 4
|
|
grpc_internal_min_log_level
|
Type: NYT::NLogging::ELogLevel
Default value: error
|
NYT::NRpc::NGrpc::TDispatcherDynamicConfig
Source: yt/yt/core/rpc/grpc/config.h
NYT::NRpc::NGrpc::TSslPemKeyCertPairConfig
Source: yt/yt/core/rpc/grpc/config.h
NYT::NRpc::TDispatcherConfig
Source: yt/yt/core/rpc/config.h
|
Parameter
|
Description
|
|
heavy_pool_size
|
Type: int
Default value: 16
|
|
compression_pool_size
|
Type: int
Default value: 8
|
|
heavy_pool_polling_period
|
Type: TDuration
Default value: 10ms
|
|
default_request_timeout
|
Type: TDuration
Default value: 1d
|
|
alert_on_missing_request_annotation
|
Type: bool
Default value: false
|
|
alert_on_unset_request_timeout
|
Type: bool
Default value: false
|
|
send_tracing_baggage
|
Type: bool
Default value: true
|
NYT::NRpc::TDispatcherDynamicConfig
Source: yt/yt/core/rpc/config.h
|
Parameter
|
Description
|
|
heavy_pool_size
|
Type: std::optional<int>
|
|
compression_pool_size
|
Type: std::optional<int>
|
|
heavy_pool_polling_period
|
Type: std::optional<TDuration>
|
|
alert_on_missing_request_annotation
|
Type: std::optional<bool>
|
|
send_tracing_baggage
|
Type: std::optional<bool>
|
NYT::NRpc::TDynamicChannelPoolConfig
Source: yt/yt/core/rpc/config.h
|
Parameter
|
Description
|
|
discover_timeout
|
Type: TDuration
Default value: 15s
|
|
acknowledgement_timeout
|
Type: TDuration
Default value: 15s
|
|
rediscover_period
|
Type: TDuration
Default value: 1m
|
|
rediscover_splay
|
Type: TDuration
Default value: 15s
|
|
hard_backoff_time
|
Type: TDuration
Default value: 1m
|
|
soft_backoff_time
|
Type: TDuration
Default value: 15s
|
|
max_peer_count
|
Type: int
Default value: 100
|
|
hashes_per_peer
|
Type: int
Default value: 10
|
|
min_peer_count_for_priority_awareness
|
Type: int
Default value: 0
|
|
enable_power_of_two_choices_strategy
|
Type: bool
Default value: true
|
|
max_concurrent_discover_requests
|
Type: int
Default value: 10
|
|
random_peer_eviction_period
|
Type: TDuration
Default value: 1s
|
|
enable_peer_polling
|
Type: bool
Default value: false
|
|
peer_polling_period
|
Type: TDuration
Default value: 1m
|
|
peer_polling_period_splay
|
Type: TDuration
Default value: 10s
|
|
peer_polling_request_timeout
|
Type: TDuration
Default value: 15s
|
|
peer_priority_strategy
|
Type: NYT::NRpc::EPeerPriorityStrategy
Default value: none
|
|
discovery_session_timeout
|
Type: TDuration
|
NYT::NRpc::THistogramExponentialBounds
Source: yt/yt/core/rpc/config.h
|
Parameter
|
Description
|
|
min
|
Type: TDuration
Default value: 0
|
|
max
|
Type: TDuration
Default value: 2s
|
NYT::NRpc::TRetryingChannelConfig
Source: yt/yt/core/rpc/config.h
|
Parameter
|
Description
|
|
retry_backoff_time
|
Type: TDuration
Default value: 3s
|
|
retry_attempts
|
Type: int
Default value: 10
|
|
enable_exponential_retry_backoffs
|
Type: bool
Default value: false
|
|
retry_backoff
|
Type: NYT::TExponentialBackoffOptions
Default value:
{
"backoff_jitter" = 0.1;
"backoff_multiplier" = 1.5;
"invocation_count" = 10;
"max_backoff" = 5000;
"min_backoff" = 1000;
}
|
|
retry_timeout
|
Type: std::optional<TDuration>
|
NYT::NRpc::TServerConfig
Source: yt/yt/core/rpc/config.h
|
Parameter
|
Description
|
|
enable_per_user_profiling
|
Type: bool
Default value: false
|
|
timing_histogram
|
Type: NYT::TIntrusivePtr<NYT::NRpc::TTimeHistogramConfig>
|
|
enable_error_code_counter
|
Type: bool
Default value: false
|
|
tracing_mode
|
Type: NYT::NRpc::ERequestTracingMode
Default value: enable
|
|
services
|
Type: THashMap<std::string, NYT::TIntrusivePtr<NYT::NYTree::INode>>
Default value: {}
|
NYT::NRpc::TServiceDiscoveryEndpointsConfig
Source: yt/yt/core/rpc/config.h
|
Parameter
|
Description
|
|
cluster
|
Type: std::optional<std::string>
|
|
clusters
|
Type: std::vector<std::string>
Default value: []
|
|
endpoint_set_id
|
Type: std::string
Required parameter
|
|
update_period
|
Type: TDuration
Default value: 1m
|
|
use_ipv4
|
Type: bool
Default value: false
|
|
use_ipv6
|
Type: bool
Default value: true
|
NYT::NRpc::TTimeHistogramConfig
Source: yt/yt/core/rpc/config.h
NYT::NTCMalloc::TDynamicHeapSizeLimitConfig
Source: yt/yt/library/tcmalloc/config.h
|
Parameter
|
Description
|
|
container_memory_ratio
|
Type: std::optional<double>
|
|
container_memory_margin
|
Type: std::optional<long>
|
|
hard
|
Type: std::optional<bool>
|
|
dump_memory_profile_on_violation
|
Type: std::optional<bool>
|
|
memory_profile_dump_timeout
|
Type: std::optional<TDuration>
|
|
memory_profile_dump_path
|
Type: std::optional<std::string>
|
|
memory_profile_dump_filename_suffix
|
Type: std::optional<std::string>
|
NYT::NTCMalloc::TDynamicTCMallocConfig
Source: yt/yt/library/tcmalloc/config.h
|
Parameter
|
Description
|
|
aggressive_release_threshold
|
Type: std::optional<long>
|
|
aggressive_release_threshold_ratio
|
Type: std::optional<double>
|
|
aggressive_release_size
|
Type: std::optional<long>
|
|
aggressive_release_period
|
Type: std::optional<TDuration>
|
|
guarded_sampling_rate
|
Type: std::optional<long>
|
|
profile_sampling_rate
|
Type: std::optional<long>
|
|
max_per_cpu_cache_size
|
Type: std::optional<long>
|
|
max_total_thread_cache_bytes
|
Type: std::optional<long>
|
|
background_release_rate
|
Type: std::optional<long>
|
|
heap_size_limit
|
Type: NYT::TIntrusivePtr<NYT::NTCMalloc::TDynamicHeapSizeLimitConfig>
Default value: {}
|
NYT::NTCMalloc::THeapSizeLimitConfig
Source: yt/yt/library/tcmalloc/config.h
|
Parameter
|
Description
|
|
container_memory_ratio
|
Type: std::optional<double>
|
|
container_memory_margin
|
Type: std::optional<long>
|
|
hard
|
Type: bool
Default value: false
|
|
dump_memory_profile_on_violation
|
Type: bool
Default value: false
|
|
memory_profile_dump_timeout
|
Type: TDuration
Default value: 10m
|
|
memory_profile_dump_path
|
Type: std::optional<std::string>
|
|
memory_profile_dump_filename_suffix
|
Type: std::optional<std::string>
|
|
memory_profile_retention
|
Type: NYT::TIntrusivePtr<NYT::NTCMalloc::TMemoryProfileRetentionConfig>
|
NYT::NTCMalloc::TMemoryProfileRetentionConfig
Source: yt/yt/library/tcmalloc/config.h
|
Parameter
|
Description
|
|
max_dump_count
|
Type: std::optional<int>
|
|
max_dump_age
|
Type: std::optional<TDuration>
|
|
max_total_size
|
Type: std::optional<long>
|
|
max_orphan_age
|
Type: TDuration
Default value: 1d
|
NYT::NTCMalloc::TTCMallocConfig
Source: yt/yt/library/tcmalloc/config.h
|
Parameter
|
Description
|
|
aggressive_release_threshold
|
Type: long
Default value: 21474836480
|
|
aggressive_release_threshold_ratio
|
Type: std::optional<double>
|
|
aggressive_release_size
|
Type: long
Default value: 134217728
|
|
aggressive_release_period
|
Type: TDuration
Default value: 100ms
|
|
guarded_sampling_rate
|
Type: std::optional<long>
Default value: 134217728
|
|
profile_sampling_rate
|
Type: long
Default value: 2097152
|
|
max_per_cpu_cache_size
|
Type: long
Default value: 3145728
|
|
max_total_thread_cache_bytes
|
Type: long
Default value: 25165824
|
|
background_release_rate
|
Type: long
Default value: 33554432
|
|
fail_fast_on_oom
|
Type: bool
Default value: true
|
|
heap_size_limit
|
Type: NYT::TIntrusivePtr<NYT::NTCMalloc::THeapSizeLimitConfig>
Default value: {}
|
NYT::NTracing::TJaegerTracerConfig
Source: yt/yt/library/tracing/jaeger/config.h
|
Parameter
|
Description
|
|
collector_channel_config
|
Type: NYT::TIntrusivePtr<NYT::NRpc::NGrpc::TChannelConfig>
|
|
flush_period
|
Type: TDuration
Default value: 15s
|
|
stop_timeout
|
Type: TDuration
Default value: 15s
|
|
rpc_timeout
|
Type: TDuration
Default value: 15s
|
|
queue_stall_timeout
|
Type: TDuration
Default value: 15m
|
|
max_request_size
|
Type: long
Default value: 131072
|
|
max_batch_size
|
Type: long
Default value: 128
|
|
max_memory
|
Type: long
Default value: 1073741824
|
|
subsampling_rate
|
Type: std::optional<double>
|
|
reconnect_period
|
Type: TDuration
Default value: 15m
|
|
endpoint_channel_timeout
|
Type: TDuration
Default value: 2h
|
|
service_name
|
Type: std::optional<std::string>
|
|
process_tags
|
Type: THashMap<std::string, std::string>
Default value: {}
|
|
enable_pid_tag
|
Type: bool
Default value: false
|
|
test_drop_spans
|
Type: bool
Default value: false
|
NYT::NTracing::TJaegerTracerDynamicConfig
Source: yt/yt/library/tracing/jaeger/config.h
|
Parameter
|
Description
|
|
collector_channel
|
Type: NYT::TIntrusivePtr<NYT::NRpc::NGrpc::TChannelConfig>
|
|
max_request_size
|
Type: std::optional<long>
|
|
max_memory
|
Type: std::optional<long>
|
|
subsampling_rate
|
Type: std::optional<double>
|
|
flush_period
|
Type: std::optional<TDuration>
|
NYT::NYTree::TSize
Description:
A wrapper type around int64 that extends its parsing from yson. It can be parsed from yson values of the following types:
-
int64. The value is used as is.
-
uint64. The value is used as is, but an exception is thrown if it exceeds the range of int64.
-
string. The value is parsed from a string with support for multiplier suffixes.
K, M, G, T, P, E — suffixes corresponding to multipliers , , etc.
Ki, Mi, Gi, Ti, Pi, Ei — suffixes corresponding to multipliers , , etc.
So 12M is , and 1Ki is .
NYT::NYson::EEnumYsonStorageType
|
Possible values
|
Description
|
|
string
|
|
|
int
|
|
NYT::NYson::EUtf8Check
|
Possible values
|
Description
|
|
disable
|
|
|
log_on_fail
|
|
|
throw_on_fail
|
|
NYT::NYson::TProtobufInteropConfig
Source: yt/yt/core/yson/config.h
|
Parameter
|
Description
|
|
default_enum_yson_storage_type
|
Type: NYT::NYson::EEnumYsonStorageType
Default value: string
|
|
utf8_check
|
Type: NYT::NYson::EUtf8Check
Default value: throw_on_fail
|
|
force_snake_case_names
|
Type: bool
Default value: false
|
|
force_enum_string_type
|
Type: bool
Default value: false
|
NYT::NYson::TProtobufInteropDynamicConfig
Source: yt/yt/core/yson/config.h
NYT::TSingletonsDynamicConfig
Source: yt/yt/core/misc/configurable_singleton_def.h
TDuration
Description:
A duration type representing a time interval. It can be parsed from yson values of the following types:
- uint64. Specifies the time in milliseconds.
- int64. Same as uint64. An exception is thrown if the value is negative.
- double. Same as int64. Behavior is undefined if the value falls outside the range of uint64.
- string. Specifies the time in formats:
10s, 15ms, 15.05s, 20us, 25. If no suffix is specified, the number is interpreted as a time in seconds.