Onyx 0.14.4
Distributed, masterless, high performance, fault tolerant data processing for Clojure
onyx.api
Public variables and functions:
- await-job-completion
- build-resume-point
- clear-checkpoints
- clear-job-data
- clear-tenancy
- gc
- gc-checkpoints
- job-ids-history
- job-snapshot-coordinates
- job-state
- kill-job
- map-set-workflow->workflow
- named-job-snapshot-coordinates
- shutdown-env
- shutdown-peer
- shutdown-peer-group
- shutdown-peers
- start-env
- start-peer-group
- start-peers
- submit-job
- subscribe-to-log
- validate-submission
onyx.checkpoint
Onyx checkpoint interfaces.
Public variables and functions:
- assume-checkpoint-coordinate
- cancel!
- complete?
- gc-checkpoint!
- gc-replica-epoch-watermark!
- read-all-replica-epoch-watermarks
- read-checkpoint
- read-checkpoint-coordinate
- read-replica-epoch-low-watermark
- stop
- storage
- watch-checkpoint-coordinate
- write-checkpoint
- write-checkpoint-coordinate
- write-replica-epoch-low-watermark
- write-replica-epoch-watermark
onyx.compression.nippy
Public variables and functions:
- checkpoint-compress
- checkpoint-compress-opts
- checkpoint-decompress
- checkpoint-decompress-opts
- messaging-compress
- messaging-compress-opts
- messaging-decompress
- messaging-decompress-opts
- statedb-compress
- statedb-decompress
- window-log-compress
- window-log-decompress
- window-log-decompress-opts
- zookeeper-compress
- zookeeper-compress-opts
- zookeeper-decompress
- zookeeper-decompress-opts
onyx.extensions
Extension interfaces for internally used queues, logs, and distributed coordination.
Public variables and functions:
- apply-log-entry
- connected?
- fire-side-effects!
- force-write-chunk
- gc-log-entry
- group-exists?
- IEmitEvent
- job-status
- monitoring-agent
- new-task-state!
- on-delete
- reactions
- read-chunk
- read-job-name-metadata
- read-log-entry
- register-pulse
- replica-diff
- subscribe-to-log
- update-origin!
- write-chunk
- write-job-name-metadata
- write-log-entry
onyx.log.commands.common
Public variables and functions:
- allocations->peers
- build-stop-task-fn
- grouped-task?
- input-task?
- job->peers
- job-allocations->peer-ids
- job-peer-count
- leaf-tasks
- messenger-slot-id
- peer->allocated-job
- promote-orphans
- remove-peers
- replica->job-peers
- root-tasks
- src-peers
- start-new-lifecycle
- start-task!
- state-task?
- stop-lifecycle-peer-group!
- to-graph
- upstream-peers
onyx.log.replica-invariants
Public variables and functions:
- active-job-invariant
- all-coordinators-exist
- all-groups-invariant
- all-jobs-have-coordinator
- all-peers-are-group-indexed
- all-peers-are-reverse-group-indexed
- all-tasks-have-non-zero-peers
- allocations-invariant
- group-index-keys-never-nil
- group-index-vals-never-nil
- no-extra-coordinators
- peer-site-invariant
- short-identifiers-correct
- slot-id-invariant
- version-invariant
onyx.log.zookeeper
Public variables and functions:
- ->ZooKeeper
- await-entry!
- catalog-path
- checkpoint-path
- checkpoint-path-version
- checkpoint-task-key
- chunk-path
- clean-up-broken-connections
- clear-job-data
- clear-tenancy
- current-replica
- epoch-low-watermark-path
- epoch-path
- exception-path
- find-log-parameters
- flow-path
- global-metadata-path
- initialize-origin!
- job-config-path
- job-hash-path
- job-metadata-path
- job-name-path
- job-paths
- latest-checkpoint-path
- lifecycles-path
- log-parameters-path
- log-path
- map->ZooKeeper
- origin-path
- pad-sequential-id
- parse-task-key
- poll-znode-timeout-ms
- prefix-path
- pulse-path
- read-checkpoint-coord
- read-log-entry
- read-log-parameters
- read-origin
- receive-event
- resume-point-path
- root-path
- savepoint-path
- seek-and-put-entry!
- seek-to-new-origin!
- task-path
- tenancy-alive?
- throw-subscriber-closed
- triggers-path
- windows-path
- workflow-path
- znode
- znode-path
- zookeeper
onyx.peer.coordinator
Public variables and functions:
- ->PeerCoordinator
- barrier-period
- check-peer-timeout!
- checkpoint-complete?
- complete-job
- completed-checkpoint
- Coordinator
- coordinator-backoff-ms
- coordinator-iteration
- emit-replica
- emit-seal-barrier
- evict-peer!
- initialise-state
- input-publications
- make-task-data
- map->PeerCoordinator
- merged-statuses
- new-peer-coordinator
- next-replica
- offer-barriers
- offer-heartbeats
- periodic-barrier
- schedule-next-barrier
- shutdown
- start-coordinator!
- start-messenger
- stop-coordinator!
- write-coordinate
onyx.peer.peer-group-manager
Public variables and functions:
- ->PeerGroupManager
- action
- annotate-reaction
- futures-stuck-park-ms
- idle-backoff-ms
- map->PeerGroupManager
- media-driver-backoff-ms
- peer-group-error-backoff-ms
- peer-group-manager
- peer-group-manager-loop
- peers-allocated-proportion
- poll-inbox!
- remove-shutdown-futs
- safe-stop-vpeer!
- setup-group-state
- shutting-down-task-metrics
- spin-park-ms
- spin-until-tasks-shutdown
- start-communicator!
- transition-group
- transition-peers
- update-scheduler-lag!
onyx.scheduling.common-job-scheduler
Public variables and functions:
- actual-usage
- add-allocation-versions
- add-messaging-short-ids
- anti-jitter-constraints
- assign-coordinators
- assign-task-resources
- assign-task-slot-ids
- btr-place-scheduling
- build-current-model
- build-job-and-task->node
- build-node->task
- build-peer->task
- build-peer->vm
- calculate-capacity
- capacity-constraints
- change-peer-allocations
- claim-spare-peers
- constrainted-tasks-for-peer
- deallocate-starved-jobs
- full-allocation?
- grouped-task?
- grouping-task-constraints
- input-task?
- job->planned-task-capacity
- job-claim-peers
- job-coverable?
- job-lower-bound
- job-offer-n-peers
- job-upper-bound
- messaging-long-form
- n-no-op-tasks
- n-qualified-peers
- no-tagged-peers-constraints
- peer-running-constraints
- reallocate-peers
- reconfigure-cluster-workload
- remove-job
- slot-ids
- src-peers
- task-tagged-constraints
- to-node-array
- unassign-task-slot-ids
- unconstrained-tasks
- unrolled-tasks
- update-peer-site
- update-slot-id-for-peer
onyx.schema
Public variables and functions:
- ->RestrictedKwNamespace
- add-event-schema
- add-state-event-schema
- AeronIdleStrategy
- BarrierCoordinate
- base-task-map
- build-allowed-key-ns
- Catalog
- combine-restricted-ns
- deprecated
- EnvConfig
- Epoch
- Event
- FlowAction
- FlowCondition
- FluxPolicy
- FnPath
- Function
- function-task-map
- GroupId
- grouping-task?
- information-model->schema
- InitialiseMode
- input-task-map
- InputResumeDefinition
- InputResumeMode
- java?
- Job
- JobConfig
- JobId
- JobMetadata
- JobName
- JobScheduler
- Language
- Lifecycle
- LifecycleCall
- LifecycleState
- LogEntry
- lookup-schema
- map->RestrictedKwNamespace
- Messaging
- NamespacedKeyword
- NonNamespacedKeyword
- output-task-map
- OutputResumeDefinition
- OutputResumeMode
- partial-clojure-fn-task
- partial-clojure-plugin
- partial-fn-task
- partial-grouping-task
- partial-input-task
- partial-java-fn-task
- partial-java-plugin
- partial-output-task
- partial-reduce-task
- PartialJob
- PartialWorkflow
- PeerClientConfig
- PeerConfig
- PeerId
- PeerSchedulerEvent
- PeerSite
- PeerZKClientConfig
- PosInt
- Reactions
- reduce-task-map
- RefinementCall
- Replica
- ReplicaDiff
- ReplicaVersion
- restricted-ns
- ResumeCoordinate
- ResumePoint
- schema-name->schema
- SegmentKey
- SlotId
- SlotMigration
- SpecialFlowTasks
- SPosInt
- State
- StateAggregationCall
- StateEvent
- StateFilterImpl
- StateLogImpl
- Storage
- TaskId
- TaskMap
- TaskName
- TaskScheduler
- TenancyId
- TenancyIdStr
- Trigger
- TriggerCall
- TriggerDelay
- TriggerDelayElements
- TriggerEvent
- TriggerEventType
- TriggerId
- TriggerPeriod
- TriggerPeriodElements
- TriggerRefinement
- TriggerState
- TriggerThreshold
- TriggerThresholdElements
- type->schema
- UniqueTaskMap
- Unit
- UnsupportedFlowKey
- UnsupportedTriggerKey
- UnsupportedWindowKey
- valid-min-peers-max-peers-n-peers?
- Window
- WindowBase
- WindowExtension
- WindowId
- WindowResumeDefinition
- WindowResumeMode
- WindowState
- WindowType
- Workflow
onyx.state.lmdb
Implementation of a LMDB backed state store. Currently this is alpha level quality.
Public variables and functions:
onyx.state.serializers.checkpoint
Streaming serializer / deserializer for state checkpoints.
Public variables and functions:
onyx.system
Public variables and functions:
- ->OnyxClient
- ->OnyxDevelopmentEnv
- ->OnyxPeer
- ->OnyxPeerGroup
- ->OnyxTask
- client-components
- development-components
- map->OnyxClient
- map->OnyxDevelopmentEnv
- map->OnyxPeer
- map->OnyxPeerGroup
- map->OnyxTask
- onyx-client
- onyx-development-env
- onyx-peer-group
- onyx-task
- onyx-vpeer-system
- peer-components
- peer-group-components
- rethrow-component
- task-components
onyx.test-helper
Public variables and functions:
- ->TestPeers
- add-test-env-peers!
- feedback-exception!
- find-task
- get-counts
- job->min-peers-per-task
- job-allocation-counts
- load-config
- map->TestPeers
- playback-log
- shutdown-env
- shutdown-peer
- shutdown-peer-group
- try-start-env
- try-start-group
- try-start-peers
- validate-enough-peers!
- with-components
- with-test-env
onyx.triggers
Public variables and functions:
- exceeds-percentile-watermark?
- exceeds-watermark?
- next-fire-time
- percentile-watermark
- percentile-watermark-fire?
- punctuation
- punctuation-fire?
- punctuation-init-locals
- punctuation-init-state
- punctuation-next-state
- segment
- segment-fire?
- segment-init-state
- segment-next-state
- timer
- timer-fire?
- timer-init-state
- timer-next-state
- watermark
- watermark-fire?
- watermark-init-locals
onyx.types
Public variables and functions:
- ->Link
- ->MonitorEvent
- ->MonitorEventBytes
- ->MonitorEventLatency
- ->MonitorEventLatencyBytes
- ->MonitorTaskEventCount
- ->Result
- ->Route
- ->StateEvent
- ->TriggerState
- barrier
- barrier-id
- checkpointed-state-event
- heartbeat
- heartbeat-id
- map->Link
- map->MonitorEvent
- map->MonitorEventBytes
- map->MonitorEventLatency
- map->MonitorEventLatencyBytes
- map->MonitorTaskEventCount
- map->Result
- map->Route
- map->StateEvent
- map->TriggerState
- max-control-message-size
- message
- message-id
- new-state-event
- ready
- ready-id
- ready-reply
- ready-reply-id
onyx.windowing.aggregation
Public variables and functions:
- average
- average-aggregation-apply-fn
- average-aggregation-create-fn
- average-aggregation-fn-init
- average-super-aggregation
- collect-by-key
- collect-by-key-aggregation-apply-log
- collect-by-key-aggregation-fn
- collect-by-key-aggregation-fn-init
- collect-by-key-super-aggregation
- collect-key-value
- collect-key-value-aggregation-apply-log
- collect-key-value-aggregation-fn
- collect-key-value-aggregation-fn-init
- collect-key-value-super-aggregation
- conj
- conj-aggregation-apply-log
- conj-aggregation-fn
- conj-aggregation-fn-init
- conj-super-aggregation
- count
- count-aggregation-apply-fn
- count-aggregation-create-state-update
- count-aggregation-fn-init
- count-super-aggregation
- max
- max-aggregation-fn
- max-create-aggregation-apply-fn
- max-super-aggregation
- min
- min-aggregation-apply-fn
- min-aggregation-create-fn
- min-super-aggregation
- set-value-aggregation-apply-log
- sum
- sum-aggregation-fn-init
- sum-aggregation-function
- sum-create-state-update-fn
- sum-super-aggregation