Poster for "On-demand Memory Compression of Stream Aggregates through Reinforcement Learning"
Abstract
This is the poster of the paper "On-demand Compression of Stream Aggregates through Reinforcement Learning", which was published at ICPE 2025.
Full text
Engineering'and' Physical'Sciences' Research'Council' Grant'number' EP/X029174/1 Horizon'Europe'2021-2027' Framework'Programme' Grant'Agreement'number'101072456 Disclaimer:+Funded+by+the+European+Union.+ Views+and+opinions+expressed+are+however+those+of+ the+author(s)+only+and+do+not+necessarily+reflect+those+ of+the+EU.+The+EU+cannot+be+held+responsible+for+them. Funded&by the&European&Union Chalmers University of Technology and University of Gothenburg, Sweden ICPE 2025 Research Track Stream Processing and Aggregates Reinforcement Learning Evalua8on Usecases and Setup SPE Controller Environment !state π "ac'on π RL Agent #reward π good ac'on posi've reward bad ac'on nega've reward - ππππ’π‘'πππ‘π - π‘βπππ’πβππ’π‘ - ππ’π‘ππ’π‘'πππ‘π - πππ‘ππππ¦ - π/π πππ‘ππ - πΆππ'ππππ π’πππ‘πππ -πππ‘ππππ¦ - π/π'πππ‘ππ - #π π‘πππ 'πππ'ππππ πππ send data A 8:00 20 A 8:03 15 F fF fF f F fF fF f Input stream Stream Processing Engine (SPE) Can run distributed/in parallel Γ spread in the Cloud-IoT con'nuum Directed Acyclic Graph Outputs π΄ππ΄,ππ, π, π !, π "##, π $%&, π '( Func'on π !π‘ ππ΄ (window advance) ππ (window size) Func'on π "## Ξ,π‘ Func'on π $%& Ξ Func'on π '( Ξ,π‘ (opt.) event &me size advance remove π output Aggregate Stream π ΓEnvironment ointerface to connect the SPE and RL Agent ΓAgent oimplement training algorithm by Neural Network (DQN) oget an ac'on to interact with the environment ΓReward ofeedback to the Agent to reinforce good ac'ons RL agent reward β΅ βΆ β· environment ΓLinear Road benchmark oVehicles travelling in highways report their posi'on/speed oEach vehicle reports its posi'on every 30 seconds oAggregate: count the number of non-consecu've stops oWS = 10 mins, WA = 5 secs Γ Synthe;c (stress-test) oData is generated following a sawtooth wave whose peaksβ values and distances are chosen randomly oThe key aYribute is generated from a Gaussian distribu'on with changing π/π) oAggregate: perform math opera'ons on a random value carried by each tuple oWS = 15 mins, WA = 1 sec ΓSetup oJava (OpenJDK 17.0.7), Python 3.7.6 oSPE: Liebre oCompression library: snappy oAgent: openAI Gym o120 episodes with maximum 1000 steps for each On-demand Memory Compression of Stream Aggregates through Reinforcement Learning Comparison discussions for the Agent with diο¬erent compression levels Scalability discussions for the Aggregate ΓWithout an Agent (top): oaverage CPU cons.: 0.33 (Linear Road), 0.59 (Synthetic) oaverage latency: 0.98s (Linear Road), 0.53s (Synthetic) ΓWith an Agent (bottom): odiff. in CPU cons. and latency are almost 0 It does not become a scaling bo5leneck for the Aggregate by introducing the Agent. Γ Linear Road (WS = 10 mins, WA = 5 secs) oall baselines are safe except for π·0.0 osimilar policy behaviors except for WEL-OB Γ Synthetic (WS = 15 mins, WA = 1 sec) oπ/π ratio decreases linearly with lower π· ofine-tune ability (ini'al) state ac'on once per day once per sec send data - πΆππππππ π ππππ β’π. π. π· β 10% -ππ‘ππ¦ β’π’ππβπππππ - πΆππππππ π πππ π β’π. π. π· β 10% oRL-based adaptive memory compression scheme for stream Aggregates oAllowing real-time balancing of performance and memory usage under latency constraints oCapture applicationand data-specific behaviors of Aggregates oHighlight the trade-off between RL training timeliness and policy effectiveness 10:00:00 ac'on 'me β’If a window hasnβt been updated for a whileβ¦ β’Compress it ο¬rstly, and later decompress it example About this paper Jingyu Liu, Vincenzo Gulisano infrequently frequently Compress! 'me for the next output 10:05:00 RL Agent SPE ac&on (compress) state, reward Cyclical dependency Γ Agent computes its next ac'on o receive the state and reward ο¬rst! baseline (X value) baseline (X value) Each π·π#baseline always sets the π· value to π β ππ#Each π·π baseline always sets the π· value to π β ππ Output Stream Γ Aggregate waits for the Agentβs ac'on o share a new state and reward ο¬rst! WEL-OB (WEL -OBlivious ) EL-OB (EL -OBlivious ) L-OB (L-OBlivious) WEL - AW (WEL-AWare) wallclock time (W) ββ β β event time (E) β β β β next output (L) (e.g. latency) β β β β -π·ππ policy observa&on Linear Road Synthetic Linear Road Synthe<c (compression threshold) π«= π β ππ, π β 0.0,1.0 : o πππππππ»πΊ βππ β₯π« Γ compress (the condition for triggering compression) β’πππ‘ππ π‘ππ: the timestamp (event time) of the latest tuple processed by the Aggregate π΄ β’π‘π : the timestamp (event time) of the latest tuple that contributed to the window instance β’e.g., 0.0 β ππ: all window instances maintained by the Aggregate π΄ are compressed 1.0 β ππ: no window instance maintained by the Aggregate π΄is compressed tuple window instance Why are four policies? ΓThree ways the system makes progress (1) wallclock 'me moves forward (2) event 'me advances (3) get a new state (e.g. latency) measurement β’(2) implies (1) because event 4me advances only as π΄ processes input, which depends on wallclock 4me. β’(3) implies (2) because latency updates occur only when event 4me advances enough to produce output. Based on these dependencies, four policies can be established: If the state comes from before 10:05:00, the eο¬ects on latency are not measured yetβ¦ TIME ALIGNMENT (in four policies) up down 100% WS π« 0% WS 100% WS π« 0% WS