1
00:00:01,240 --> 00:00:04,420
Usually we mine data that sits
somewhere in a database, or

2
00:00:04,420 --> 00:00:05,820
a distributed file system.

3
00:00:07,220 --> 00:00:12,070
And we can access the same data
repeatedly, and it is all available to us,

4
00:00:12,070 --> 00:00:12,890
whenever we need it.

5
00:00:14,112 --> 00:00:17,850
But there's some applications where data
doesn't really live in a database or

6
00:00:17,850 --> 00:00:18,850
if it does.

7
00:00:18,850 --> 00:00:19,690
The data base is so

8
00:00:19,690 --> 00:00:23,140
large, that we can't query it fast
enough to answer questions about it.

9
00:00:24,290 --> 00:00:26,580
Examples include click stream
at that are major internet site,

10
00:00:26,580 --> 00:00:29,660
at a major internet site or observational
data coming down from satellites.

11
00:00:30,782 --> 00:00:34,842
Answering queries about this sort of data
requires clever observation techniques and

12
00:00:34,842 --> 00:00:35,398
methods for

13
00:00:35,398 --> 00:00:39,250
compressing data, in a way that allows us
to answer the queries we need to answer.

14
00:00:41,710 --> 00:00:45,360
To begin with a brief summery,
of a stream management system,

15
00:00:45,360 --> 00:00:50,140
the analog of the data base management
system, the idea of sliding windows is

16
00:00:50,140 --> 00:00:55,800
an essential idea that tells, that lets us
focus on, on recent data in the stream,

17
00:00:55,800 --> 00:00:58,910
we'll then discuss a particular problem,
that of counting ones and

18
00:00:58,910 --> 00:01:06,083
the windows of a big stream gets
slide by The fundamental difference

19
00:01:06,083 --> 00:01:10,260
between a data stream and a database, is
who controls how data enters the system.

20
00:01:12,760 --> 00:01:15,630
In a database system, the staff
associated with the management of

21
00:01:15,630 --> 00:01:19,440
the database generally insert data
into the system using a bulk loader or

22
00:01:19,440 --> 00:01:22,340
even explicit SQL INSERT commands.

23
00:01:22,340 --> 00:01:26,590
The staff can decide how much data to
load into the system, when, and how fast.

24
00:01:27,650 --> 00:01:32,410
In a streaming environment, the management
cannot control the rate of input.

25
00:01:32,410 --> 00:01:35,340
For example, the search queries that
arrive at Google are generated by

26
00:01:35,340 --> 00:01:38,480
random people around the world,
at their pace.

27
00:01:38,480 --> 00:01:41,450
Google staff have no control
over the arrival of queries.

28
00:01:41,450 --> 00:01:45,040
They have to architect their system,
to deal with whatever data rate there is.

29
00:01:47,350 --> 00:01:51,050
You might think a transaction processing
system like Walmart recording all

30
00:01:51,050 --> 00:01:55,400
the purchases at all its cash
registers everywhere, as a stream.

31
00:01:55,400 --> 00:01:57,260
And in a sense it is.

32
00:01:57,260 --> 00:02:01,330
But Walmart has a large but
fixed number of registers.

33
00:02:01,330 --> 00:02:04,820
And checkout clerks can
press the keys just so fast.

34
00:02:04,820 --> 00:02:08,380
So there's actually a pretty well defined
limit on how fast data arrives in

35
00:02:08,380 --> 00:02:09,185
such a system.

36
00:02:09,185 --> 00:02:16,451
[SOUND] So let's see the elements of
the data stream model of computation.

37
00:02:16,451 --> 00:02:20,070
First we assume inputs are tuples
just as in the database system.

38
00:02:20,070 --> 00:02:23,420
Although in many algorithms we
shall assume input elements.

39
00:02:23,420 --> 00:02:27,110
Our tuples is a very simple
kind like bits or, or integers.

40
00:02:28,270 --> 00:02:32,250
We assume there, are one or
more input ports to which data arrives.

41
00:02:32,250 --> 00:02:36,450
Generally we assume the arrival rate is
high although we will be a little vague on

42
00:02:36,450 --> 00:02:37,760
how high is high.

43
00:02:39,630 --> 00:02:44,150
The important property of the arrival
rate is that it is fast enough.

44
00:02:44,150 --> 00:02:48,480
That it is not feasible for the system
to store all the arriving data and

45
00:02:48,480 --> 00:02:51,250
at the same time make it
instantaneously available for

46
00:02:51,250 --> 00:02:56,410
any query we might want
to perform on the data.

47
00:02:56,410 --> 00:03:00,380
As a result, the interesting algorithms or
data stream which are generally methods

48
00:03:00,380 --> 00:03:04,800
that use a limited amount of storage,
perhaps only main memory.

49
00:03:04,800 --> 00:03:08,410
And still enable us to answer important
queries about the content to the stream.

50
00:03:13,390 --> 00:03:16,470
Streams can be queried in two modes.

51
00:03:16,470 --> 00:03:20,300
The first is similar to the way
we query a database system.

52
00:03:20,300 --> 00:03:21,570
You ask a query once and

53
00:03:21,570 --> 00:03:25,410
expect an answer about the state of
the system at the time you ask the query.

54
00:03:27,710 --> 00:03:30,730
For example,
what is the maximum value seen in

55
00:03:30,730 --> 00:03:34,540
the screen from its beginning to
the exact time the query is asked?

56
00:03:34,540 --> 00:03:37,200
This question can be answered
by keeping a single value,

57
00:03:37,200 --> 00:03:44,975
the maximum, and updating it if necessary
each time a new stream element arrives.

58
00:03:44,975 --> 00:03:45,510
'Kay.

59
00:03:45,510 --> 00:03:47,850
The other kind of query is
called a standing query.

60
00:03:47,850 --> 00:03:49,840
You write the query once.

61
00:03:49,840 --> 00:03:53,830
And you expect the system to make
the answer available at all times,

62
00:03:53,830 --> 00:03:57,660
perhaps outputting a new value
each time the answer changes.

63
00:03:57,660 --> 00:04:02,310
For instance, a standing query might ask
for a report of each stream element that

64
00:04:02,310 --> 00:04:05,880
is larger than any element seen so
far in the stream.

65
00:04:05,880 --> 00:04:09,310
We can answer this one by
keeping one value, the maximum.

66
00:04:09,310 --> 00:04:13,810
And each new element is compared with the
max and if it is larger we do two things.

67
00:04:13,810 --> 00:04:17,820
We output the value and
we update the max to be that value.

68
00:04:21,800 --> 00:04:26,735
So here is a very simple outline what
a stream management system looks like.

69
00:04:29,761 --> 00:04:30,570
Kay.

70
00:04:30,570 --> 00:04:36,940
First there's a processor which is
the software, that executes the queries.

71
00:04:36,940 --> 00:04:41,380
The processor, could of course be a large
number of processors working in concert.

72
00:04:41,380 --> 00:04:45,339
The processor may store
some standing queries.

73
00:04:47,380 --> 00:04:50,420
And also allow ad-hoc queries
to be issued by the user.

74
00:04:57,200 --> 00:04:59,430
Here we see several streams
entering the system.

75
00:05:02,330 --> 00:05:05,120
Conditionally we'll assume that
the element at the right end of

76
00:05:05,120 --> 00:05:08,300
the stream has arrived most recently.

77
00:05:08,300 --> 00:05:10,560
And time goes backward to the left.

78
00:05:10,560 --> 00:05:13,818
That is the further, left the earlier
the element entered the system.

79
00:05:13,818 --> 00:05:19,556
[SOUND] The system makes
outputs in response,

80
00:05:19,556 --> 00:05:25,905
to the standing queries and
the add hock queries.

81
00:05:25,905 --> 00:05:28,310
Now usually there is
some archival storage.

82
00:05:28,310 --> 00:05:29,310
This storage is so

83
00:05:29,310 --> 00:05:33,550
massive, that it is not possible to
do more than store the input streams.

84
00:05:33,550 --> 00:05:37,990
We cannot assume the archival storage
is architected like a database system,

85
00:05:37,990 --> 00:05:39,800
where by using appropriate indices or

86
00:05:39,800 --> 00:05:44,270
other tools, one can answer queries
efficiently from that data.

87
00:05:44,270 --> 00:05:48,160
We only know that if we had to reconstruct
the history of the streams we could.

88
00:05:48,160 --> 00:05:50,200
Perhaps taking a long time to do so.

89
00:05:52,540 --> 00:05:57,710
Now, there is a limit at working storage,
which might be main memory, flash storage,

90
00:05:57,710 --> 00:05:58,960
or even disk.

91
00:05:58,960 --> 00:06:01,860
But we assume it holds essential
parts of the input streams in

92
00:06:01,860 --> 00:06:03,840
a way that supports fast querying.

93
00:06:10,820 --> 00:06:15,090
We're going to list some examples
of the sorts of streams that it

94
00:06:15,090 --> 00:06:16,070
could be useful to mine.

95
00:06:17,420 --> 00:06:21,000
One example is the query stream
at a search engine like Google.

96
00:06:21,000 --> 00:06:22,370
For example, Google Trends,

97
00:06:22,370 --> 00:06:26,130
wants to find out which search queries are
much more frequent today, than yesterday.

98
00:06:28,130 --> 00:06:31,060
These queries represent issues
of rising public interest.

99
00:06:32,540 --> 00:06:36,090
Answering such a standing query requires
looking back at most two days in

100
00:06:36,090 --> 00:06:37,320
the query stream.

101
00:06:37,320 --> 00:06:40,140
That's quite a lot,
perhaps billions of queries.

102
00:06:40,140 --> 00:06:43,789
But it is tiny compared with the stream
of all Google queries ever issued.

103
00:06:48,040 --> 00:06:50,510
Click streams are another
source of very rapid input.

104
00:06:51,650 --> 00:06:54,620
A site like yahoo has many
millions of users each day,

105
00:06:54,620 --> 00:06:57,480
and the average user probably
clicks a dozen times or more.

106
00:06:58,640 --> 00:07:01,870
A question worth answering is
which URLs are getting clicked on,

107
00:07:01,870 --> 00:07:03,839
a lot more this past hour than normally.

108
00:07:04,970 --> 00:07:07,230
Interestingly, while some
of these events refler,

109
00:07:07,230 --> 00:07:11,940
reflect breaking news stories,
many also represent a broken link.

110
00:07:11,940 --> 00:07:14,220
When people can't get the page they want,

111
00:07:14,220 --> 00:07:17,330
they'll often click on it
several times before giving up.

112
00:07:17,330 --> 00:07:20,579
So sites mine their click streams to,
to detect broken links.

113
00:07:24,950 --> 00:07:27,010
We can view a switch in
the middle of the Internet,

114
00:07:27,010 --> 00:07:29,480
as processing streams,
one stream for each port.

115
00:07:30,810 --> 00:07:34,220
The elements of the stream,
are IP packets, typically.

116
00:07:34,220 --> 00:07:36,960
And, the switch can store a lot
of information about packets,

117
00:07:36,960 --> 00:07:39,660
including the response speed
of different network links and

118
00:07:39,660 --> 00:07:42,810
the points of origin and
destination of the packets.

119
00:07:42,810 --> 00:07:46,440
This information could be used to advice
the switch in the best routing for

120
00:07:46,440 --> 00:07:49,530
a packet, or
to detect a denial of service attack.

121
00:07:54,850 --> 00:07:57,140
Now the concept of the sliding window,
is essential for

122
00:07:57,140 --> 00:07:59,489
many of the algorithms
we're going to discuss.

123
00:08:01,910 --> 00:08:05,310
The simplest form of window,
is defined by a fixed length and

124
00:08:05,310 --> 00:08:09,880
consistent of the, most recent and
elementary received on a stream.

125
00:08:09,880 --> 00:08:11,850
Notice that each time
an element is received,

126
00:08:11,850 --> 00:08:13,970
the oldest element falls
out of the window.

127
00:08:16,660 --> 00:08:20,630
A variation is to define the window as
all the elements that have arrived within

128
00:08:20,630 --> 00:08:23,250
some time interval T
extending into the path.

129
00:08:23,250 --> 00:08:24,810
Say the last hour.

130
00:08:24,810 --> 00:08:27,840
This sort of window has a storage
requirement that is not fixed since

131
00:08:27,840 --> 00:08:30,920
the number of arrivals
within time t can vary.

132
00:08:30,920 --> 00:08:34,290
In comparison defining the window to
have a fixed number of elements lets us

133
00:08:34,290 --> 00:08:37,570
rely on needing storage space,
only up to a certain limit.

134
00:08:37,570 --> 00:08:42,400
The interesting case is when

135
00:08:42,400 --> 00:08:45,430
we're using a window consisting
of the last end string element.

136
00:08:45,430 --> 00:08:48,560
But N is so large that we cannot
store any elements in main memory.

137
00:08:50,920 --> 00:08:54,849
And while we have options to recr,
increase, the, the size of main memory.

138
00:08:55,920 --> 00:09:00,890
Use, many compute notes to handle one
window or use disc in some cases.

139
00:09:00,890 --> 00:09:03,730
We also need to consider the case,
where there are many streams,

140
00:09:03,730 --> 00:09:08,160
perhaps millions of streams,
arriving at the same stream processor.

141
00:09:08,160 --> 00:09:10,910
In that case, N does not have
to be very large, before we

142
00:09:10,910 --> 00:09:15,880
cannot store all the windows in a way that
allows us to get exact answers to queries.

143
00:09:15,880 --> 00:09:17,540
About the contents of the windows.

144
00:09:22,110 --> 00:09:25,370
So here's a little picture of a steam and
a window of length six.

145
00:09:25,370 --> 00:09:31,830
Okay initially,
the stream has arrived up to this point J.

146
00:09:33,790 --> 00:09:37,159
The elements K L and so
on will arrive in the future.

147
00:09:41,180 --> 00:09:41,900
Okay.

148
00:09:41,900 --> 00:09:46,500
Now k arrives, the oldest element s,
is no longer part of

149
00:09:46,500 --> 00:09:50,090
the window which continues to hold
exactly six elements, as it always will.

150
00:09:53,770 --> 00:09:58,830
Now l arrives and d falls out of
the the window, and z arrives,

151
00:10:00,400 --> 00:10:02,165
causing f to be dropped from the window.

152
00:10:02,165 --> 00:10:10,463
[SOUND] Let's take
a really simple example.

153
00:10:10,463 --> 00:10:12,380
Okay, we have a stream of integers.

154
00:10:13,420 --> 00:10:15,680
The window is of size N.

155
00:10:15,680 --> 00:10:18,970
That is, the window will hold the N,
most recent integers in the stream.

156
00:10:20,970 --> 00:10:24,390
And we want the system to be able
to answer one standing quero,

157
00:10:24,390 --> 00:10:27,970
query, what is the average of
the elements in the window.

158
00:10:29,790 --> 00:10:33,040
Often we imagine stream extends
infinitely into the past.

159
00:10:33,040 --> 00:10:36,580
So, we don't worry about what happen
before there had been enough arrivals to

160
00:10:36,580 --> 00:10:37,990
fill the window.

161
00:10:37,990 --> 00:10:40,880
However realistically we have
to get started some how.

162
00:10:40,880 --> 00:10:43,230
So, lets store the first N
inputs as they arrives and

163
00:10:43,230 --> 00:10:48,630
maintain the sum accountive elements seen
so far, until the account reaches the end.

164
00:10:48,630 --> 00:10:51,814
The average is the sum divided
by the count at any point.

165
00:10:55,138 --> 00:10:55,638
now.

166
00:10:57,690 --> 00:11:01,440
Suppose we have our window full, and it
consists of the most recent end elements.

167
00:11:02,620 --> 00:11:05,450
Those will store the average of
these elements that averages in

168
00:11:05,450 --> 00:11:07,830
the local storage but
it's not part of the window.

169
00:11:09,080 --> 00:11:10,580
Suppose new element I arrives.

170
00:11:11,680 --> 00:11:15,640
The oldest element J in the window
will fall out, of the window.

171
00:11:17,000 --> 00:11:22,760
Thus, the change in the average
is i minus j, all divided by N.

172
00:11:24,030 --> 00:11:29,130
I over n accounts for the contribution
i makes to the average, and minus j

173
00:11:29,130 --> 00:11:33,709
over N accounts for the fact that j no
longer equ, contributes to the average.

174
00:11:34,730 --> 00:11:38,840
The important point is that, in this
matter, we can answer the query what is

175
00:11:38,840 --> 00:11:42,990
the average of the elements in the window,
doing only a small fixed number of

176
00:11:42,990 --> 00:11:45,850
arithmetic steps,
with each arrival on the stream.

177
00:11:45,850 --> 00:11:48,700
That is far, far better than
having to compute the sum and

178
00:11:48,700 --> 00:11:53,350
average of all N elements in the window
each time a new element arrives.

179
00:11:53,350 --> 00:11:56,240
But not every query about
the current value of

180
00:11:56,240 --> 00:11:59,370
the window can be answered in
an equally convenient way.

